Team Ai
Apppublic

oex2003/evolution-api

sourceHugging Faceupdated 1y agoView on Hugging Face
0likes
base-chatbot.service.ts413 linesDownload Raw Back to chatbot
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