oex2003/evolution-api
0
1import { InstanceDto } from '@api/dto/instance.dto';2import { PrismaRepository } from '@api/repository/repository.service';3import { WAMonitoringService } from '@api/services/monitor.service';4import { Integration } from '@api/types/wa.types';5import { ConfigService } from '@config/env.config';6import { Logger } from '@config/logger.config';7import { IntegrationSession } from '@prisma/client';8 9/**10 * Base class for all chatbot service implementations11 * Contains common methods shared across different chatbot integrations12 */13export abstract class BaseChatbotService<BotType = any, SettingsType = any> {14 protected readonly logger: Logger;15 protected readonly waMonitor: WAMonitoringService;16 protected readonly prismaRepository: PrismaRepository;17 protected readonly configService?: ConfigService;18 19 constructor(20 waMonitor: WAMonitoringService,21 prismaRepository: PrismaRepository,22 loggerName: string,23 configService?: ConfigService,24 ) {25 this.waMonitor = waMonitor;26 this.prismaRepository = prismaRepository;27 this.logger = new Logger(loggerName);28 this.configService = configService;29 }30 31 /**32 * Check if a message contains an image33 */34 protected isImageMessage(content: string): boolean {35 return content.includes('imageMessage');36 }37 38 /**39 * Check if a message contains audio40 */41 protected isAudioMessage(content: string): boolean {42 return content.includes('audioMessage');43 }44 45 /**46 * Check if a string is valid JSON47 */48 protected isJSON(str: string): boolean {49 try {50 JSON.parse(str);51 return true;52 } catch (e) {53 return false;54 }55 }56 57 /**58 * Determine the media type from a URL based on its extension59 */60 protected getMediaType(url: string): string | null {61 const extension = url.split('.').pop()?.toLowerCase();62 const imageExtensions = ['jpg', 'jpeg', 'png', 'gif', 'bmp', 'webp'];63 const audioExtensions = ['mp3', 'wav', 'aac', 'ogg'];64 const videoExtensions = ['mp4', 'avi', 'mkv', 'mov'];65 const documentExtensions = ['pdf', 'doc', 'docx', 'xls', 'xlsx', 'ppt', 'pptx', 'txt'];66 67 if (imageExtensions.includes(extension || '')) return 'image';68 if (audioExtensions.includes(extension || '')) return 'audio';69 if (videoExtensions.includes(extension || '')) return 'video';70 if (documentExtensions.includes(extension || '')) return 'document';71 return null;72 }73 74 /**75 * Create a new chatbot session76 */77 public async createNewSession(instance: InstanceDto | any, data: any, type: string) {78 try {79 // Extract pushName safely - if data.pushName is an object with a pushName property, use that80 const pushNameValue =81 typeof data.pushName === 'object' && data.pushName?.pushName82 ? data.pushName.pushName83 : typeof data.pushName === 'string'84 ? data.pushName85 : null;86 87 // Extract remoteJid safely88 const remoteJidValue =89 typeof data.remoteJid === 'object' && data.remoteJid?.remoteJid ? data.remoteJid.remoteJid : data.remoteJid;90 91 const session = await this.prismaRepository.integrationSession.create({92 data: {93 remoteJid: remoteJidValue,94 pushName: pushNameValue,95 sessionId: remoteJidValue,96 status: 'opened',97 awaitUser: false,98 botId: data.botId,99 instanceId: instance.instanceId,100 type: type,101 },102 });103 104 return { session };105 } catch (error) {106 this.logger.error(error);107 return;108 }109 }110 111 /**112 * Standard implementation for processing incoming messages113 * This handles the common workflow across all chatbot types:114 * 1. Check for existing session or create new one115 * 2. Handle message based on session state116 */117 public async process(118 instance: any,119 remoteJid: string,120 bot: BotType,121 session: IntegrationSession,122 settings: SettingsType,123 content: string,124 pushName?: string,125 msg?: any,126 ): Promise<void> {127 try {128 // For new sessions or sessions awaiting initialization129 if (!session) {130 await this.initNewSession(instance, remoteJid, bot, settings, session, content, pushName, msg);131 return;132 }133 134 // If session is paused, ignore the message135 if (session.status === 'paused') {136 return;137 }138 139 // For existing sessions, keywords might indicate the conversation should end140 const keywordFinish = (settings as any)?.keywordFinish || '';141 const normalizedContent = content.toLowerCase().trim();142 if (keywordFinish.length > 0 && normalizedContent === keywordFinish.toLowerCase()) {143 // Update session to closed and return144 await this.prismaRepository.integrationSession.update({145 where: {146 id: session.id,147 },148 data: {149 status: 'closed',150 },151 });152 return;153 }154 155 // Forward the message to the chatbot API156 await this.sendMessageToBot(instance, session, settings, bot, remoteJid, pushName || '', content, msg);157 158 // Update session to indicate we're waiting for user response159 await this.prismaRepository.integrationSession.update({160 where: {161 id: session.id,162 },163 data: {164 status: 'opened',165 awaitUser: true,166 },167 });168 } catch (error) {169 this.logger.error(`Error in process: ${error}`);170 return;171 }172 }173 174 /**175 * Standard implementation for sending messages to WhatsApp176 * This handles common patterns like markdown links and formatting177 */178 protected async sendMessageWhatsApp(179 instance: any,180 remoteJid: string,181 message: string,182 settings: SettingsType,183 ): Promise<void> {184 if (!message) return;185 186 const linkRegex = /!?\[(.*?)\]\((.*?)\)/g;187 let textBuffer = '';188 let lastIndex = 0;189 let match: RegExpExecArray | null;190 191 const splitMessages = (settings as any)?.splitMessages ?? false;192 193 while ((match = linkRegex.exec(message)) !== null) {194 const [fullMatch, altText, url] = match;195 const mediaType = this.getMediaType(url);196 const beforeText = message.slice(lastIndex, match.index);197 198 if (beforeText) {199 textBuffer += beforeText;200 }201 202 if (mediaType) {203 // Send accumulated text before sending media204 if (textBuffer.trim()) {205 await this.sendFormattedText(instance, remoteJid, textBuffer.trim(), settings, splitMessages);206 textBuffer = '';207 }208 209 // Handle sending the media210 try {211 if (mediaType === 'audio') {212 await instance.audioWhatsapp({213 number: remoteJid.split('@')[0],214 delay: (settings as any)?.delayMessage || 1000,215 audio: url,216 caption: altText,217 });218 } else {219 await instance.mediaMessage(220 {221 number: remoteJid.split('@')[0],222 delay: (settings as any)?.delayMessage || 1000,223 mediatype: mediaType,224 media: url,225 caption: altText,226 fileName: mediaType === 'document' ? altText || 'document' : undefined,227 },228 null,229 false,230 );231 }232 } catch (error) {233 this.logger.error(`Error sending media: ${error}`);234 // If media fails, at least send the alt text and URL235 textBuffer += `${altText}: ${url}`;236 }237 } else {238 // It's a regular link, keep it in the text239 textBuffer += fullMatch;240 }241 242 lastIndex = linkRegex.lastIndex;243 }244 245 // Add any remaining text after the last match246 if (lastIndex < message.length) {247 const remainingText = message.slice(lastIndex);248 if (remainingText.trim()) {249 textBuffer += remainingText;250 }251 }252 253 // Send any remaining text254 if (textBuffer.trim()) {255 await this.sendFormattedText(instance, remoteJid, textBuffer.trim(), settings, splitMessages);256 }257 }258 259 /**260 * Helper method to send formatted text with proper typing indicators and delays261 */262 private async sendFormattedText(263 instance: any,264 remoteJid: string,265 text: string,266 settings: any,267 splitMessages: boolean,268 ): Promise<void> {269 const timePerChar = settings?.timePerChar ?? 0;270 const minDelay = 1000;271 const maxDelay = 20000;272 273 if (splitMessages) {274 const multipleMessages = text.split('\n\n');275 for (let index = 0; index < multipleMessages.length; index++) {276 const message = multipleMessages[index];277 if (!message.trim()) continue;278 279 const delay = Math.min(Math.max(message.length * timePerChar, minDelay), maxDelay);280 281 if (instance.integration === Integration.WHATSAPP_BAILEYS) {282 await instance.client.presenceSubscribe(remoteJid);283 await instance.client.sendPresenceUpdate('composing', remoteJid);284 }285 286 await new Promise<void>((resolve) => {287 setTimeout(async () => {288 await instance.textMessage(289 {290 number: remoteJid.split('@')[0],291 delay: settings?.delayMessage || 1000,292 text: message,293 },294 false,295 );296 resolve();297 }, delay);298 });299 300 if (instance.integration === Integration.WHATSAPP_BAILEYS) {301 await instance.client.sendPresenceUpdate('paused', remoteJid);302 }303 }304 } else {305 const delay = Math.min(Math.max(text.length * timePerChar, minDelay), maxDelay);306 307 if (instance.integration === Integration.WHATSAPP_BAILEYS) {308 await instance.client.presenceSubscribe(remoteJid);309 await instance.client.sendPresenceUpdate('composing', remoteJid);310 }311 312 await new Promise<void>((resolve) => {313 setTimeout(async () => {314 await instance.textMessage(315 {316 number: remoteJid.split('@')[0],317 delay: settings?.delayMessage || 1000,318 text: text,319 },320 false,321 );322 resolve();323 }, delay);324 });325 326 if (instance.integration === Integration.WHATSAPP_BAILEYS) {327 await instance.client.sendPresenceUpdate('paused', remoteJid);328 }329 }330 }331 332 /**333 * Standard implementation for initializing a new session334 * This method should be overridden if a subclass needs specific initialization335 */336 protected async initNewSession(337 instance: any,338 remoteJid: string,339 bot: BotType,340 settings: SettingsType,341 session: IntegrationSession,342 content: string,343 pushName?: string | any,344 msg?: any,345 ): Promise<void> {346 // Create a session if none exists347 if (!session) {348 // Extract pushName properly - if it's an object with pushName property, use that349 const pushNameValue =350 typeof pushName === 'object' && pushName?.pushName351 ? pushName.pushName352 : typeof pushName === 'string'353 ? pushName354 : null;355 356 const sessionResult = await this.createNewSession(357 {358 instanceName: instance.instanceName,359 instanceId: instance.instanceId,360 },361 {362 remoteJid,363 pushName: pushNameValue,364 botId: (bot as any).id,365 },366 this.getBotType(),367 );368 369 if (!sessionResult || !sessionResult.session) {370 this.logger.error('Failed to create new session');371 return;372 }373 374 session = sessionResult.session;375 }376 377 // Update session status to opened378 await this.prismaRepository.integrationSession.update({379 where: {380 id: session.id,381 },382 data: {383 status: 'opened',384 awaitUser: false,385 },386 });387 388 // Forward the message to the chatbot389 await this.sendMessageToBot(instance, session, settings, bot, remoteJid, pushName || '', content, msg);390 }391 392 /**393 * Get the bot type identifier (e.g., 'dify', 'n8n', 'evoai')394 * This should match the type field used in the IntegrationSession395 */396 protected abstract getBotType(): string;397 398 /**399 * Send a message to the chatbot API400 * This is specific to each chatbot integration401 */402 protected abstract sendMessageToBot(403 instance: any,404 session: IntegrationSession,405 settings: SettingsType,406 bot: BotType,407 remoteJid: string,408 pushName: string,409 content: string,410 msg?: any,411 ): Promise<void>;412}413 