Team Ai
Apppublic

oex2003/evolution-api

sourceHugging Faceupdated 1y agoView on Hugging Face
0likes
openai.service.ts715 linesDownload Raw Back to services
1import { PrismaRepository } from '@api/repository/repository.service';2import { WAMonitoringService } from '@api/services/monitor.service';3import { Integration } from '@api/types/wa.types';4import { ConfigService, Language, Openai as OpenaiConfig } from '@config/env.config';5import { IntegrationSession, OpenaiBot, OpenaiSetting } from '@prisma/client';6import { sendTelemetry } from '@utils/sendTelemetry';7import axios from 'axios';8import { downloadMediaMessage } from 'baileys';9import FormData from 'form-data';10import OpenAI from 'openai';11import P from 'pino';12 13import { BaseChatbotService } from '../../base-chatbot.service';14 15/**16 * OpenAI service that extends the common BaseChatbotService17 * Handles both Assistant API and ChatCompletion API18 */19export class OpenaiService extends BaseChatbotService<OpenaiBot, OpenaiSetting> {20  protected client: OpenAI;21 22  constructor(waMonitor: WAMonitoringService, prismaRepository: PrismaRepository, configService: ConfigService) {23    super(waMonitor, prismaRepository, 'OpenaiService', configService);24  }25 26  /**27   * Return the bot type for OpenAI28   */29  protected getBotType(): string {30    return 'openai';31  }32 33  /**34   * Initialize the OpenAI client with the provided API key35   */36  protected initClient(apiKey: string) {37    this.client = new OpenAI({ apiKey });38    return this.client;39  }40 41  /**42   * Process a message based on the bot type (assistant or chat completion)43   */44  public async process(45    instance: any,46    remoteJid: string,47    openaiBot: OpenaiBot,48    session: IntegrationSession,49    settings: OpenaiSetting,50    content: string,51    pushName?: string,52    msg?: any,53  ): Promise<void> {54    try {55      this.logger.log(`Starting process for remoteJid: ${remoteJid}, bot type: ${openaiBot.botType}`);56 57      // Handle audio message transcription58      if (content.startsWith('audioMessage|') && msg) {59        this.logger.log('Detected audio message, attempting to transcribe');60 61        // Get OpenAI credentials for transcription62        const creds = await this.prismaRepository.openaiCreds.findUnique({63          where: { id: openaiBot.openaiCredsId },64        });65 66        if (!creds) {67          this.logger.error(`OpenAI credentials not found. CredsId: ${openaiBot.openaiCredsId}`);68          return;69        }70 71        // Initialize OpenAI client for transcription72        this.initClient(creds.apiKey);73 74        // Transcribe the audio75        const transcription = await this.speechToText(msg, instance);76 77        if (transcription) {78          this.logger.log(`Audio transcribed: ${transcription}`);79          // Replace the audio message identifier with the transcription80          content = transcription;81        } else {82          this.logger.error('Failed to transcribe audio');83          await this.sendMessageWhatsApp(84            instance,85            remoteJid,86            "Sorry, I couldn't transcribe your audio message. Could you please type your message instead?",87            settings,88          );89          return;90        }91      } else {92        // Get the OpenAI credentials93        const creds = await this.prismaRepository.openaiCreds.findUnique({94          where: { id: openaiBot.openaiCredsId },95        });96 97        if (!creds) {98          this.logger.error(`OpenAI credentials not found. CredsId: ${openaiBot.openaiCredsId}`);99          return;100        }101 102        // Initialize OpenAI client103        this.initClient(creds.apiKey);104      }105 106      // Handle keyword finish107      const keywordFinish = settings?.keywordFinish || '';108      const normalizedContent = content.toLowerCase().trim();109      if (keywordFinish.length > 0 && normalizedContent === keywordFinish.toLowerCase()) {110        if (settings?.keepOpen) {111          await this.prismaRepository.integrationSession.update({112            where: {113              id: session.id,114            },115            data: {116              status: 'closed',117            },118          });119        } else {120          await this.prismaRepository.integrationSession.delete({121            where: {122              id: session.id,123            },124          });125        }126 127        await sendTelemetry('/openai/session/finish');128        return;129      }130 131      // If session is new or doesn't exist132      if (!session) {133        const data = {134          remoteJid,135          pushName,136          botId: openaiBot.id,137        };138 139        const createSession = await this.createNewSession(140          { instanceName: instance.instanceName, instanceId: instance.instanceId },141          data,142          this.getBotType(),143        );144 145        await this.initNewSession(146          instance,147          remoteJid,148          openaiBot,149          settings,150          createSession.session,151          content,152          pushName,153          msg,154        );155 156        await sendTelemetry('/openai/session/start');157        return;158      }159 160      // If session exists but is paused161      if (session.status === 'paused') {162        await this.prismaRepository.integrationSession.update({163          where: {164            id: session.id,165          },166          data: {167            status: 'opened',168            awaitUser: true,169          },170        });171 172        return;173      }174 175      // Process with the appropriate API based on bot type176      await this.sendMessageToBot(instance, session, settings, openaiBot, remoteJid, pushName || '', content);177    } catch (error) {178      this.logger.error(`Error in process: ${error.message || JSON.stringify(error)}`);179      return;180    }181  }182 183  /**184   * Send message to OpenAI - this handles both Assistant API and ChatCompletion API185   */186  protected async sendMessageToBot(187    instance: any,188    session: IntegrationSession,189    settings: OpenaiSetting,190    openaiBot: OpenaiBot,191    remoteJid: string,192    pushName: string,193    content: string,194  ): Promise<void> {195    this.logger.log(`Sending message to bot for remoteJid: ${remoteJid}, bot type: ${openaiBot.botType}`);196 197    if (!this.client) {198      this.logger.log('Client not initialized, initializing now');199      const creds = await this.prismaRepository.openaiCreds.findUnique({200        where: { id: openaiBot.openaiCredsId },201      });202 203      if (!creds) {204        this.logger.error(`OpenAI credentials not found in sendMessageToBot. CredsId: ${openaiBot.openaiCredsId}`);205        return;206      }207 208      this.initClient(creds.apiKey);209    }210 211    try {212      let message: string;213 214      // Handle different bot types215      if (openaiBot.botType === 'assistant') {216        this.logger.log('Processing with Assistant API');217        message = await this.processAssistantMessage(218          instance,219          session,220          openaiBot,221          remoteJid,222          pushName,223          false, // Not fromMe224          content,225        );226      } else {227        this.logger.log('Processing with ChatCompletion API');228        message = await this.processChatCompletionMessage(instance, openaiBot, remoteJid, content);229      }230 231      this.logger.log(`Got response from OpenAI: ${message?.substring(0, 50)}${message?.length > 50 ? '...' : ''}`);232 233      // Send the response234      if (message) {235        this.logger.log('Sending message to WhatsApp');236        await this.sendMessageWhatsApp(instance, remoteJid, message, settings);237      } else {238        this.logger.error('No message to send to WhatsApp');239      }240 241      // Update session status242      await this.prismaRepository.integrationSession.update({243        where: {244          id: session.id,245        },246        data: {247          status: 'opened',248          awaitUser: true,249        },250      });251    } catch (error) {252      this.logger.error(`Error in sendMessageToBot: ${error.message || JSON.stringify(error)}`);253      if (error.response) {254        this.logger.error(`API Response data: ${JSON.stringify(error.response.data || {})}`);255      }256      return;257    }258  }259 260  /**261   * Process message using the OpenAI Assistant API262   */263  private async processAssistantMessage(264    instance: any,265    session: IntegrationSession,266    openaiBot: OpenaiBot,267    remoteJid: string,268    pushName: string,269    fromMe: boolean,270    content: string,271  ): Promise<string> {272    const messageData: any = {273      role: fromMe ? 'assistant' : 'user',274      content: [{ type: 'text', text: content }],275    };276 277    // Handle image messages278    if (this.isImageMessage(content)) {279      const contentSplit = content.split('|');280      const url = contentSplit[1].split('?')[0];281 282      messageData.content = [283        { type: 'text', text: contentSplit[2] || content },284        {285          type: 'image_url',286          image_url: {287            url: url,288          },289        },290      ];291    }292 293    // Get thread ID from session or create new thread294    let threadId = session.sessionId;295 296    // Create a new thread if one doesn't exist or invalid format297    if (!threadId || threadId === remoteJid) {298      const newThread = await this.client.beta.threads.create();299      threadId = newThread.id;300 301      // Save the new thread ID to the session302      await this.prismaRepository.integrationSession.update({303        where: {304          id: session.id,305        },306        data: {307          sessionId: threadId,308        },309      });310      this.logger.log(`Created new thread ID: ${threadId} for session: ${session.id}`);311    }312 313    // Add message to thread314    await this.client.beta.threads.messages.create(threadId, messageData);315 316    if (fromMe) {317      sendTelemetry('/message/sendText');318      return '';319    }320 321    // Run the assistant322    const runAssistant = await this.client.beta.threads.runs.create(threadId, {323      assistant_id: openaiBot.assistantId,324    });325 326    if (instance.integration === Integration.WHATSAPP_BAILEYS) {327      await instance.client.presenceSubscribe(remoteJid);328      await instance.client.sendPresenceUpdate('composing', remoteJid);329    }330 331    // Wait for the assistant to complete332    const response = await this.getAIResponse(threadId, runAssistant.id, openaiBot.functionUrl, remoteJid, pushName);333 334    if (instance.integration === Integration.WHATSAPP_BAILEYS) {335      await instance.client.sendPresenceUpdate('paused', remoteJid);336    }337 338    // Extract the response text safely with type checking339    let responseText = "I couldn't generate a proper response. Please try again.";340    try {341      const messages = response?.data || [];342      if (messages.length > 0) {343        const messageContent = messages[0]?.content || [];344        if (messageContent.length > 0) {345          const textContent = messageContent[0];346          if (textContent && 'text' in textContent && textContent.text && 'value' in textContent.text) {347            responseText = textContent.text.value;348          }349        }350      }351    } catch (error) {352      this.logger.error(`Error extracting response text: ${error}`);353    }354 355    // Update session with the thread ID to ensure continuity356    await this.prismaRepository.integrationSession.update({357      where: {358        id: session.id,359      },360      data: {361        status: 'opened',362        awaitUser: true,363        sessionId: threadId, // Ensure thread ID is saved consistently364      },365    });366 367    // Return fallback message if unable to extract text368    return responseText;369  }370 371  /**372   * Process message using the OpenAI ChatCompletion API373   */374  private async processChatCompletionMessage(375    instance: any,376    openaiBot: OpenaiBot,377    remoteJid: string,378    content: string,379  ): Promise<string> {380    this.logger.log('Starting processChatCompletionMessage');381 382    // Check if client is initialized383    if (!this.client) {384      this.logger.log('Client not initialized in processChatCompletionMessage, initializing now');385      const creds = await this.prismaRepository.openaiCreds.findUnique({386        where: { id: openaiBot.openaiCredsId },387      });388 389      if (!creds) {390        this.logger.error(`OpenAI credentials not found. CredsId: ${openaiBot.openaiCredsId}`);391        return 'Error: OpenAI credentials not found';392      }393 394      this.initClient(creds.apiKey);395    }396 397    // Check if model is defined398    if (!openaiBot.model) {399      this.logger.error('OpenAI model not defined');400      return 'Error: OpenAI model not configured';401    }402 403    this.logger.log(`Using model: ${openaiBot.model}, max tokens: ${openaiBot.maxTokens || 500}`);404 405    // Get existing conversation history from the session406    const session = await this.prismaRepository.integrationSession.findFirst({407      where: {408        remoteJid,409        botId: openaiBot.id,410        status: 'opened',411      },412    });413 414    let conversationHistory = [];415 416    if (session && session.context) {417      try {418        const sessionData =419          typeof session.context === 'string' ? JSON.parse(session.context as string) : session.context;420 421        conversationHistory = sessionData.history || [];422        this.logger.log(`Retrieved conversation history from session, ${conversationHistory.length} messages`);423      } catch (error) {424        this.logger.error(`Error parsing session context: ${error.message}`);425        // Continue with empty history if we can't parse the session data426        conversationHistory = [];427      }428    }429 430    // Log bot data431    this.logger.log(`Bot data - systemMessages: ${JSON.stringify(openaiBot.systemMessages || [])}`);432    this.logger.log(`Bot data - assistantMessages: ${JSON.stringify(openaiBot.assistantMessages || [])}`);433    this.logger.log(`Bot data - userMessages: ${JSON.stringify(openaiBot.userMessages || [])}`);434 435    // Prepare system messages436    const systemMessages: any = openaiBot.systemMessages || [];437    const messagesSystem: any[] = systemMessages.map((message) => {438      return {439        role: 'system',440        content: message,441      };442    });443 444    // Prepare assistant messages445    const assistantMessages: any = openaiBot.assistantMessages || [];446    const messagesAssistant: any[] = assistantMessages.map((message) => {447      return {448        role: 'assistant',449        content: message,450      };451    });452 453    // Prepare user messages454    const userMessages: any = openaiBot.userMessages || [];455    const messagesUser: any[] = userMessages.map((message) => {456      return {457        role: 'user',458        content: message,459      };460    });461 462    // Prepare current message463    const messageData: any = {464      role: 'user',465      content: [{ type: 'text', text: content }],466    };467 468    // Handle image messages469    if (this.isImageMessage(content)) {470      this.logger.log('Found image message');471      const contentSplit = content.split('|');472      const url = contentSplit[1].split('?')[0];473 474      messageData.content = [475        { type: 'text', text: contentSplit[2] || content },476        {477          type: 'image_url',478          image_url: {479            url: url,480          },481        },482      ];483    }484 485    // Combine all messages: system messages, pre-defined messages, conversation history, and current message486    const messages: any[] = [487      ...messagesSystem,488      ...messagesAssistant,489      ...messagesUser,490      ...conversationHistory,491      messageData,492    ];493 494    this.logger.log(`Final messages payload: ${JSON.stringify(messages)}`);495 496    if (instance.integration === Integration.WHATSAPP_BAILEYS) {497      this.logger.log('Setting typing indicator');498      await instance.client.presenceSubscribe(remoteJid);499      await instance.client.sendPresenceUpdate('composing', remoteJid);500    }501 502    // Send the request to OpenAI503    try {504      this.logger.log('Sending request to OpenAI API');505      const completions = await this.client.chat.completions.create({506        model: openaiBot.model,507        messages: messages,508        max_tokens: openaiBot.maxTokens || 500, // Add default if maxTokens is missing509      });510 511      if (instance.integration === Integration.WHATSAPP_BAILEYS) {512        await instance.client.sendPresenceUpdate('paused', remoteJid);513      }514 515      const responseContent = completions.choices[0].message.content;516      this.logger.log(`Received response from OpenAI: ${JSON.stringify(completions.choices[0])}`);517 518      // Add the current exchange to the conversation history and update the session519      conversationHistory.push(messageData);520      conversationHistory.push({521        role: 'assistant',522        content: responseContent,523      });524 525      // Limit history length to avoid token limits (keep last 10 messages)526      if (conversationHistory.length > 10) {527        conversationHistory = conversationHistory.slice(conversationHistory.length - 10);528      }529 530      // Save the updated conversation history to the session531      if (session) {532        await this.prismaRepository.integrationSession.update({533          where: { id: session.id },534          data: {535            context: JSON.stringify({536              history: conversationHistory,537            }),538          },539        });540        this.logger.log(`Updated session with conversation history, now ${conversationHistory.length} messages`);541      }542 543      return responseContent;544    } catch (error) {545      this.logger.error(`Error calling OpenAI: ${error.message || JSON.stringify(error)}`);546      if (error.response) {547        this.logger.error(`API Response status: ${error.response.status}`);548        this.logger.error(`API Response data: ${JSON.stringify(error.response.data || {})}`);549      }550      return `Sorry, there was an error: ${error.message || 'Unknown error'}`;551    }552  }553 554  /**555   * Wait for and retrieve the AI response556   */557  private async getAIResponse(558    threadId: string,559    runId: string,560    functionUrl: string | null,561    remoteJid: string,562    pushName: string,563  ) {564    let status = await this.client.beta.threads.runs.retrieve(threadId, runId);565 566    let maxRetries = 60; // 1 minute with 1s intervals567    const checkInterval = 1000; // 1 second568 569    while (570      status.status !== 'completed' &&571      status.status !== 'failed' &&572      status.status !== 'cancelled' &&573      status.status !== 'expired' &&574      maxRetries > 0575    ) {576      await new Promise((resolve) => setTimeout(resolve, checkInterval));577      status = await this.client.beta.threads.runs.retrieve(threadId, runId);578 579      // Handle tool calls580      if (status.status === 'requires_action' && status.required_action?.type === 'submit_tool_outputs') {581        const toolCalls = status.required_action.submit_tool_outputs.tool_calls;582        const toolOutputs = [];583 584        for (const toolCall of toolCalls) {585          if (functionUrl) {586            try {587              const payloadData = JSON.parse(toolCall.function.arguments);588 589              // Add context590              payloadData.remoteJid = remoteJid;591              payloadData.pushName = pushName;592 593              const response = await axios.post(functionUrl, {594                functionName: toolCall.function.name,595                functionArguments: payloadData,596              });597 598              toolOutputs.push({599                tool_call_id: toolCall.id,600                output: JSON.stringify(response.data),601              });602            } catch (error) {603              this.logger.error(`Error calling function: ${error}`);604              toolOutputs.push({605                tool_call_id: toolCall.id,606                output: JSON.stringify({ error: 'Function call failed' }),607              });608            }609          } else {610            toolOutputs.push({611              tool_call_id: toolCall.id,612              output: JSON.stringify({ error: 'No function URL configured' }),613            });614          }615        }616 617        await this.client.beta.threads.runs.submitToolOutputs(threadId, runId, {618          tool_outputs: toolOutputs,619        });620      }621 622      maxRetries--;623    }624 625    if (status.status === 'completed') {626      const messages = await this.client.beta.threads.messages.list(threadId);627      return messages;628    } else {629      this.logger.error(`Assistant run failed with status: ${status.status}`);630      return { data: [{ content: [{ text: { value: 'Failed to get a response from the assistant.' } }] }] };631    }632  }633 634  protected isImageMessage(content: string): boolean {635    return content.includes('imageMessage');636  }637 638  /**639   * Implementation of speech-to-text transcription for audio messages640   */641  public async speechToText(msg: any, instance: any): Promise<string | null> {642    const settings = await this.prismaRepository.openaiSetting.findFirst({643      where: {644        instanceId: instance.instanceId,645      },646    });647 648    if (!settings) {649      this.logger.error(`OpenAI settings not found. InstanceId: ${instance.instanceId}`);650      return null;651    }652 653    const creds = await this.prismaRepository.openaiCreds.findUnique({654      where: { id: settings.openaiCredsId },655    });656 657    if (!creds) {658      this.logger.error(`OpenAI credentials not found. CredsId: ${settings.openaiCredsId}`);659      return null;660    }661 662    let audio: Buffer;663 664    if (msg.message.mediaUrl) {665      audio = await axios.get(msg.message.mediaUrl, { responseType: 'arraybuffer' }).then((response) => {666        return Buffer.from(response.data, 'binary');667      });668    } else if (msg.message.base64) {669      audio = Buffer.from(msg.message.base64, 'base64');670    } else {671      // Fallback for raw WhatsApp audio messages that need downloadMediaMessage672      audio = await downloadMediaMessage(673        { key: msg.key, message: msg?.message },674        'buffer',675        {},676        {677          logger: P({678        customLevels: {679          verbose: 15,680          debug: 20,681          info: 30,682          warn: 40,683          error: 50,684          fatal: 60,685        },686        level: 'error',687        useOnlyCustomLevels: true,688      }) as any,689          reuploadRequest: instance,690        },691      );692    }693 694    const lang = this.configService.get<Language>('LANGUAGE').includes('pt')695      ? 'pt'696      : this.configService.get<Language>('LANGUAGE');697 698    const formData = new FormData();699    formData.append('file', audio, 'audio.ogg');700    formData.append('model', 'whisper-1');701    formData.append('language', lang);702 703    const apiKey = creds?.apiKey || this.configService.get<OpenaiConfig>('OPENAI').API_KEY_GLOBAL;704 705    const response = await axios.post('https://api.openai.com/v1/audio/transcriptions', formData, {706      headers: {707        'Content-Type': 'multipart/form-data',708        Authorization: `Bearer ${apiKey}`,709      },710    });711 712    return response?.data?.text;713  }714}715