import { StaticAudienceExecutor } from './executors/StaticAudienceExecutor' import { CriteriaParser } from './builders/CriteriaParser' import { SupabaseContactRepository } from './repositories/SupabaseContactRepository' import { AudienceCriteria, AudienceQueryOptions, AudienceExecutionResult, AudienceBuilderConfig, CreateAudienceData, UpdateAudienceData } from './types' /** * Módulo centralizado para todas as operações de Audience * * @example * ```typescript * const audienceModule = new AudienceModule() * * // Configurar repositórios necessários * audienceModule.setRepositories({ * contactRepository, * audienceRepository, * memberRepository * }) * * // Executar query de audiência * const result = await audienceModule.executeQuery(criteria, { * organizationId: 'org-123', * projectId: 'proj-456' * }) * ``` */ export class AudienceModule { private audienceRepository: any private memberRepository: any private contactRepository: any private staticExecutor: StaticAudienceExecutor constructor(config: AudienceBuilderConfig = {}) { this.staticExecutor = new StaticAudienceExecutor() // Se supabaseClient foi fornecido, o módulo vira "self-contained" para queries (engine V2 no próprio módulo) if (config.supabaseClient) { const repo = new SupabaseContactRepository({ supabaseClient: config.supabaseClient, clickhouseClient: config.clickhouseClient, debug: !!config.enableDebugLogs, logger: console }) this.setRepositories({ contactRepository: repo }) } // Log config se debug estiver habilitado if (config.enableDebugLogs) { // ⚠️ Nunca logar `config` inteiro aqui: pode conter `supabaseClient` (com URL/chaves) e outros objetos sensíveis. console.log('[AUDIENCE_MODULE] configurado', { enableDebugLogs: true, hasSupabaseClient: !!config.supabaseClient }) } } /** * Configura os repositórios necessários * Deve ser chamado antes de usar o módulo */ setRepositories(repositories: { contactRepository?: any audienceRepository?: any memberRepository?: any }) { if (repositories.contactRepository) { this.contactRepository = repositories.contactRepository this.staticExecutor.setContactRepository(repositories.contactRepository) } if (repositories.audienceRepository) { this.audienceRepository = repositories.audienceRepository } if (repositories.memberRepository) { this.memberRepository = repositories.memberRepository } } // =========================== // QUERY EXECUTION // =========================== /** * Executa query de audiência e retorna resultados */ async executeQuery( criteria: string | AudienceCriteria, options: AudienceQueryOptions ): Promise { const parsed = CriteriaParser.parse(criteria) const type = CriteriaParser.getAudienceType(parsed) if (type === 'static') { return this.staticExecutor.execute(parsed, options) } else { return this.executeLiveQuery(parsed, options) } } /** * Retorna apenas a contagem de contatos (mais rápido) */ async getContactCount( criteria: string | AudienceCriteria, options: Omit ): Promise { const parsed = CriteriaParser.parse(criteria) return this.staticExecutor.executeCount(parsed, options) } /** * Retorna IDs de contatos que atendem aos critérios */ async getContactIds( criteria: string | AudienceCriteria, organizationId: string, projectId: string ): Promise> { if (!this.contactRepository) { throw new Error('ContactRepository não configurado') } const parsed = CriteriaParser.parse(criteria) return this.contactRepository.getContactIdsByAudienceCriteriaV2( organizationId, projectId, parsed ) } /** * Avalia se um contato específico atende aos critérios (ótimo para Journeys/Live) * sem precisar calcular o conjunto completo de IDs. */ async matchesContact( criteria: string | AudienceCriteria, organizationId: string, projectId: string, contactId: string, options?: { triggerEvent?: { event_id?: string; event_name?: string; event_timestamp?: string } } ): Promise { if (!this.contactRepository) { throw new Error('ContactRepository não configurado') } const parsed = CriteriaParser.parse(criteria) if (typeof this.contactRepository.matchesContactByAudienceCriteriaV2 === 'function') { return this.contactRepository.matchesContactByAudienceCriteriaV2( organizationId, projectId, parsed, contactId, options?.triggerEvent ) } // Fallback (mais lento): calcular ids e testar presença const ids = await this.contactRepository.getContactIdsByAudienceCriteriaV2( organizationId, projectId, parsed ) return ids.has(contactId) } // =========================== // CRUD OPERATIONS // =========================== /** * Cria nova audiência */ async createAudience(data: CreateAudienceData) { if (!this.audienceRepository) { throw new Error('AudienceRepository não configurado') } const validation = CriteriaParser.validate(data.criteria) if (!validation.valid) { throw new Error(`Critérios inválidos: ${validation.errors.join(', ')}`) } const { data: audience, error } = await this.audienceRepository.create(data) if (error) throw error const type = CriteriaParser.getAudienceType(data.criteria) if (type === 'static') { const count = await this.getContactCount(data.criteria, { organizationId: data.organization_id, projectId: data.project_id }) await this.audienceRepository.updateCount(audience!.id, count) audience!.count = count } return { data: audience, error: null } } /** * Atualiza audiência existente */ async updateAudience( id: string, updateData: UpdateAudienceData, organizationId: string, projectId: string ) { if (!this.audienceRepository) { throw new Error('AudienceRepository não configurado') } if (updateData.criteria) { const validation = CriteriaParser.validate(updateData.criteria) if (!validation.valid) { throw new Error(`Critérios inválidos: ${validation.errors.join(', ')}`) } } const { data: audience, error } = await this.audienceRepository.update(id, updateData) if (error) throw error if (updateData.criteria) { const type = CriteriaParser.getAudienceType(updateData.criteria) if (type === 'static') { const count = await this.getContactCount(updateData.criteria, { organizationId, projectId }) await this.audienceRepository.updateCount(id, count) audience!.count = count } } return { data: audience, error: null } } /** * Busca audiência por ID */ async getAudienceById(id: string, organizationId: string, projectId: string) { if (!this.audienceRepository) { throw new Error('AudienceRepository não configurado') } return this.audienceRepository.findById(id, organizationId, projectId) } /** * Deleta audiência */ async deleteAudience(id: string) { if (!this.audienceRepository) { throw new Error('AudienceRepository não configurado') } return this.audienceRepository.delete(id) } // =========================== // MEMBERS MANAGEMENT // =========================== /** * Verifica se uma audience tem regra de evento específico */ hasEventRule(criteria: any, eventName: string): boolean { try { // Parse do critério let parsed = criteria if (typeof criteria === 'string') { try { parsed = JSON.parse(criteria) } catch { return false } } // Verificar formato V2 com groups/rules if (parsed.groups && Array.isArray(parsed.groups)) { for (const group of parsed.groups) { if (group.rules && Array.isArray(group.rules)) { const found = group.rules.some((rule: any) => rule.kind === 'event' && rule.eventName === eventName ) if (found) return true } } } return false } catch (error) { return false } } /** * Adiciona contatos à audiência (bulk) */ async addMembers( audienceId: string, contactIds: string[], organizationId: string, projectId: string, origin: 'realtime' | 'backfill' = 'backfill' ) { if (!this.memberRepository) { throw new Error('MemberRepository não configurado') } return this.memberRepository.bulkUpsert( audienceId, organizationId, projectId, contactIds, origin ) } /** * Lista membros da audiência */ async getMembers( audienceId: string, organizationId: string, projectId: string, pagination?: { page?: number; limit?: number } ) { if (!this.memberRepository) { throw new Error('MemberRepository não configurado') } return this.memberRepository.listMembers( audienceId, organizationId, projectId, pagination ) } // =========================== // UTILITIES // =========================== /** * Valida critérios de audiência */ validateCriteria(criteria: string | AudienceCriteria) { const parsed = CriteriaParser.parse(criteria) return CriteriaParser.validate(parsed) } /** * Detecta tipo de audiência */ getAudienceType(criteria: string | AudienceCriteria): 'static' | 'live' { const parsed = CriteriaParser.parse(criteria) return CriteriaParser.getAudienceType(parsed) } // =========================== // PRIVATE METHODS // =========================== /** * Executa query de audiência live */ private async executeLiveQuery( _criteria: AudienceCriteria, options: AudienceQueryOptions ): Promise { const startTime = Date.now() if (!this.memberRepository) { throw new Error('MemberRepository não configurado para audiences live') } const { data: members, count } = await this.memberRepository.listMembers( '', options.organizationId, options.projectId, options.pagination ) const contactIds = new Set(members?.map((m: any) => m.id) || []) return { contactIds, contacts: members || [], count: count || 0, metadata: { executionTime: Date.now() - startTime, criteriaType: 'live' } } } }