oex2003/evolution-api
0
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 