diff --git a/.env.example b/.env.example index ab11a13..9f835c6 100644 --- a/.env.example +++ b/.env.example @@ -6,7 +6,15 @@ DEFAULT_MODEL=openai-responses RUNNER_KEY=R2D2C3PO # Conversation History (Generic) -CONVERSATION_HISTORY_MAX_MESSAGES=60 +CONVERSATION_HISTORY_TTL=2629800 + +# Conversation Cache Global switch +CONVERSATION_CACHE_ENABLED=true + +# Optional provider-specific overrides +OPENAI_CONVERSATION_CACHE_ENABLED=true +GEMINI_CONVERSATION_CACHE_ENABLED=true +# ... # Ollama Configuration OLLAMA_BASE_URL=http://localhost:11434 diff --git a/README.md b/README.md index 96d7c5a..eae5420 100644 --- a/README.md +++ b/README.md @@ -10,6 +10,23 @@ API for interacting with LEIA instances. ## Usage +### Environment variables + +Use `.env.example` as reference. For conversation cache behavior in Redis, you can control it with one switch: + +- Global: `CONVERSATION_CACHE_ENABLED=true|false` +- TTL de conversaciones: `CONVERSATION_HISTORY_TTL=2629800` (seconds; default 2629800 = 1 month) +- Provider overrides: + - `OPENAI_CONVERSATION_CACHE_ENABLED=true|false` + - `GEMINI_CONVERSATION_CACHE_ENABLED=true|false` + - `OLLAMA_CONVERSATION_CACHE_ENABLED=true|false` + +Priority order: + +1. Provider-specific variable +2. Global variable +3. Default value (`true`) + ### Start the server ```bash diff --git a/controllers/evaluationController.js b/controllers/evaluationController.js index f78f185..e459292 100644 --- a/controllers/evaluationController.js +++ b/controllers/evaluationController.js @@ -2,6 +2,8 @@ const modelManager = require('../models/modelManager'); const sessionService = require('../services/sessionService'); module.exports.evaluateSolution = async function evaluateSolution(req, res) { + let conversationClearable = false; + try { const { sessionId, result } = req.body; @@ -21,6 +23,8 @@ module.exports.evaluateSolution = async function evaluateSolution(req, res) { return res.status(404).send({ error: `LEIA metadata for session ID: ${sessionId} not found` }); } + conversationClereable = true; + // Obtener el modelo const model = modelManager.getModel(sessionData.modelName); @@ -34,5 +38,13 @@ module.exports.evaluateSolution = async function evaluateSolution(req, res) { } catch (error) { console.error(`Error evaluating solution for session ${req.body.sessionId}:`, error); res.status(500).send({ error: 'Internal error evaluating solution' }); + } finally { + if (conversationClereable) { + try { + await sessionService.clearConversation(req.body.sessionId); + } catch (cleanupError) { + console.error(`Error clearing conversation cache for session ${req.body.sessionId}:`, cleanupError); + } + } } }; \ No newline at end of file diff --git a/models/conversationStore.js b/models/conversationStore.js deleted file mode 100644 index e9006b3..0000000 --- a/models/conversationStore.js +++ /dev/null @@ -1,201 +0,0 @@ -const { redisClient } = require('../config/redis'); - -/** - * ConversationStore manages conversation history for any provider that requires context management. - * Uses Redis to persist and retrieve conversation messages, supporting message normalization, - * system instruction management, and automatic history trimming. - */ -class ConversationStore { - /** - * Creates a new ConversationStore instance - * @param {Object} options - Configuration options - * @param {string} options.prefix - Redis key prefix (default: 'session:conversation:') - * @param {string} options.providerName - Provider name for env var lookup (default: generic settings) - * @param {number} options.defaultMaxMessages - Default max messages when not configured (default: 60) - */ - constructor(options = {}) { - this.keyPrefix = options.prefix || 'conversation:'; - this.providerName = options.providerName || ''; - this.defaultMaxMessages = options.defaultMaxMessages || 60; - this.maxMessages = this.parseMaxMessages(); - } - - /** - * Parses max messages from environment variables - * Checks provider-specific env var first, then falls back to generic one - * @private - * @returns {number} Maximum messages to keep in history - */ - parseMaxMessages() { - let rawValue; - - if (this.providerName) { - const providerEnvVar = `${this.providerName.toUpperCase()}_HISTORY_MAX_MESSAGES`; - rawValue = process.env[providerEnvVar]; - } - - if (!rawValue) { - rawValue = process.env.CONVERSATION_HISTORY_MAX_MESSAGES; - } - - const parsed = Number.parseInt(rawValue || this.defaultMaxMessages, 10); - - if (!Number.isInteger(parsed) || parsed <= 0) { - return this.defaultMaxMessages; - } - - return parsed; - } - - getConversationKey(sessionId) { - return `${this.keyPrefix}${sessionId}`; - } - - /** - * Normalizes and validates a message - * @param {string} role - Message role (system, user, assistant) - * @param {string} content - Message content - * @returns {Object|null} Normalized message or null if invalid - */ - normalizeMessage(role, content) { - const normalizedRole = typeof role === 'string' ? role.trim() : ''; - const normalizedContent = typeof content === 'string' ? content.trim() : ''; - - if (!['system', 'user', 'assistant'].includes(normalizedRole)) { - return null; - } - - if (!normalizedContent) { - return null; - } - - return { - role: normalizedRole, - content: normalizedContent, - }; - } - - /** - * Retrieves the full conversation history for a session - * @param {string} sessionId - Session identifier - * @returns {Promise} Array of normalized messages - */ - async getConversation(sessionId) { - const rawMessages = await redisClient.lRange(this.getConversationKey(sessionId), 0, -1); - - return rawMessages - .map((rawMessage) => { - try { - return JSON.parse(rawMessage); - } catch (error) { - return null; - } - }) - .map((message) => (message ? this.normalizeMessage(message.role, message.content) : null)) - .filter(Boolean); - } - - /** - * Appends a message to the conversation history and trims if necessary - * @param {string} sessionId - Session identifier - * @param {string} role - Message role - * @param {string} content - Message content - * @returns {Promise} - */ - async appendMessage(sessionId, role, content) { - const message = this.normalizeMessage(role, content); - - if (!message) { - return; - } - - const key = this.getConversationKey(sessionId); - await redisClient.rPush(key, JSON.stringify(message)); - await redisClient.lTrim(key, -this.maxMessages, -1); - } - - /** - * Ensures a system message exists at the start of conversation - * Creates one if missing, updates if present, or adds as first message if not already there - * @param {string} sessionId - Session identifier - * @param {string} systemInstruction - System instruction text - * @returns {Promise} - */ - async ensureSystemMessage(sessionId, systemInstruction) { - const normalizedSystemMessage = this.normalizeMessage('system', systemInstruction); - - if (!normalizedSystemMessage) { - return; - } - - const key = this.getConversationKey(sessionId); - const firstRawMessage = await redisClient.lIndex(key, 0); - - if (!firstRawMessage) { - await redisClient.rPush(key, JSON.stringify(normalizedSystemMessage)); - return; - } - - let firstMessage = null; - - try { - firstMessage = JSON.parse(firstRawMessage); - } catch (error) { - firstMessage = null; - } - - const normalizedFirstMessage = firstMessage - ? this.normalizeMessage(firstMessage.role, firstMessage.content) - : null; - - if (!normalizedFirstMessage) { - await redisClient.lSet(key, 0, JSON.stringify(normalizedSystemMessage)); - return; - } - - if (normalizedFirstMessage.role !== 'system') { - await redisClient.lPush(key, JSON.stringify(normalizedSystemMessage)); - await redisClient.lTrim(key, -this.maxMessages, -1); - return; - } - - if (normalizedFirstMessage.content !== normalizedSystemMessage.content) { - await redisClient.lSet(key, 0, JSON.stringify(normalizedSystemMessage)); - } - } - - /** - * Builds a complete conversation for a provider request - * Ensures system message, adds user message, and returns full history - * @param {string} sessionId - Session identifier - * @param {string} systemInstruction - System instruction to ensure - * @param {string} userMessage - User message to append - * @returns {Promise} Complete conversation history ready for LLM - */ - async buildConversationForRequest(sessionId, systemInstruction, userMessage) { - await this.ensureSystemMessage(sessionId, systemInstruction); - await this.appendMessage(sessionId, 'user', userMessage); - return this.getConversation(sessionId); - } - - /** - * Stores a provider response in the conversation history - * @param {string} sessionId - Session identifier - * @param {string} assistantMessage - Assistant/model response to store - * @returns {Promise} - */ - async storeAssistantResponse(sessionId, assistantMessage) { - await this.appendMessage(sessionId, 'assistant', assistantMessage); - } - - /** - * Completely clears the conversation history for a session - * @param {string} sessionId - Session identifier - * @returns {Promise} - */ - async clearConversation(sessionId) { - await redisClient.del(this.getConversationKey(sessionId)); - } -} - -module.exports.ConversationStore = ConversationStore; \ No newline at end of file diff --git a/models/providers/baseModel.js b/models/providers/baseModel.js index 8be6bdb..4a60669 100644 --- a/models/providers/baseModel.js +++ b/models/providers/baseModel.js @@ -1,12 +1,202 @@ const Errors = require('../../utils/errors'); const Prompts = require('../../utils/prompts'); const ProviderState = require('../providerState'); +const Environment = require('../../utils/environment'); +const { redisClient } = require('../../config/redis'); class BaseModel { constructor() { this.name = 'base'; - this.apiKeyEnvVar = ''; + this.native = true; + this.envVar = null; this._client = null; + this.conversationPrefix = 'conversations:'; + this.defaultConversationTtlSeconds = 2629800; + } + + // conversationStore methods implemented for all providers by default + + /** + * Obtiene la clave Redis para una conversación de sesión. + * @param {string} sessionId - Session ID + * @returns {string} + */ + getConversationKey(sessionId) { + const normalizedSessionId = typeof sessionId === 'string' ? sessionId.trim() : ''; + return `${this.conversationPrefix}${normalizedSessionId}`; + } + + /** + * Obtiene TTL de historial de conversación en segundos. + * @returns {number} + */ + getConversationTtlSeconds() { + const rawTtlSeconds = + process.env.CONVERSATION_HISTORY_TTL ?? + process.env.CONVERSATION_HISTORY_TTL_SECONDS ?? + this.defaultConversationTtlSeconds; + const parsed = Number.parseInt(rawTtlSeconds, 10); + if (!Number.isInteger(parsed) || parsed <= 0) { + return this.defaultConversationTtlSeconds; + } + + return parsed; + } + + /** + * Normaliza un mensaje de conversación. + * @param {string} role - Rol del mensaje + * @param {string} content - Contenido del mensaje + * @returns {Object|null} + */ + normalizeConversationMessage(role, content) { + const normalizedRole = typeof role === 'string' ? role.trim() : ''; + const normalizedContent = typeof content === 'string' ? content.trim() : ''; + + if (!['system', 'user', 'assistant'].includes(normalizedRole)) { + return null; + } + + if (!normalizedContent) { + return null; + } + + return { + role: normalizedRole, + content: normalizedContent, + }; + } + + /** + * Obtiene el historial completo de la conversación. + * @param {string} sessionId - Session ID + * @returns {Promise} + */ + async getConversation(sessionId) { + if (!Environment.isCacheEnabled(this.envVar, this.native)) { + return []; + } + + const rawMessages = await redisClient.lRange(this.getConversationKey(sessionId), 0, -1); + + return rawMessages + .map((rawMessage) => { + try { + return JSON.parse(rawMessage); + } catch (error) { + return null; + } + }) + .map((message) => (message ? this.normalizeConversationMessage(message.role, message.content) : null)) + .filter(Boolean); + } + + /** + * Añade un mensaje al historial y recorta al máximo permitido. + * @param {string} sessionId - Session ID + * @param {string} role - Rol del mensaje + * @param {string} content - Contenido del mensaje + * @returns {Promise} + */ + async appendMessage(sessionId, role, content) { + const message = this.normalizeConversationMessage(role, content); + + if (!message || !Environment.isCacheEnabled(this.envVar, this.native)) { + return; + } + + const key = this.getConversationKey(sessionId); + await redisClient.rPush(key, JSON.stringify(message)); + await redisClient.expire(key, this.getConversationTtlSeconds()); + } + + /** + * Garantiza que exista un mensaje de sistema al inicio del historial. + * @param {string} sessionId - Session ID + * @param {string} systemInstruction - Instrucción de sistema + * @returns {Promise} + */ + async ensureSystemMessage(sessionId, systemInstruction) { + const normalizedSystemMessage = this.normalizeConversationMessage('system', systemInstruction); + + if (!normalizedSystemMessage || !Environment.isCacheEnabled(this.envVar, this.native)) { + return; + } + + const key = this.getConversationKey(sessionId); + const firstRawMessage = await redisClient.lIndex(key, 0); + + if (!firstRawMessage) { + await redisClient.rPush(key, JSON.stringify(normalizedSystemMessage)); + await redisClient.expire(key, this.getConversationTtlSeconds()); + return; + } + + let firstMessage = null; + + try { + firstMessage = JSON.parse(firstRawMessage); + } catch (error) { + firstMessage = null; + } + + const normalizedFirstMessage = firstMessage + ? this.normalizeConversationMessage(firstMessage.role, firstMessage.content) + : null; + + if (!normalizedFirstMessage) { + await redisClient.lSet(key, 0, JSON.stringify(normalizedSystemMessage)); + await redisClient.expire(key, this.getConversationTtlSeconds()); + return; + } + + if (normalizedFirstMessage.role !== 'system') { + await redisClient.lPush(key, JSON.stringify(normalizedSystemMessage)); + await redisClient.expire(key, this.getConversationTtlSeconds()); + return; + } + + if (normalizedFirstMessage.content !== normalizedSystemMessage.content) { + await redisClient.lSet(key, 0, JSON.stringify(normalizedSystemMessage)); + } + + await redisClient.expire(key, this.getConversationTtlSeconds()); + } + + /** + * Construye la conversación completa para una solicitud al proveedor. + * @param {string} sessionId - Session ID + * @param {string} systemInstruction - Instrucción de sistema + * @param {string} userMessage - Mensaje del usuario + * @returns {Promise} + */ + async buildConversationForRequest(sessionId, systemInstruction, userMessage) { + await this.ensureSystemMessage(sessionId, systemInstruction); + await this.appendMessage(sessionId, 'user', userMessage); + return this.getConversation(sessionId); + } + + /** + * Guarda la respuesta del asistente en el historial. + * @param {string} sessionId - Session ID + * @param {string} assistantMessage - Mensaje del asistente + * @returns {Promise} + */ + async storeAssistantResponse(sessionId, assistantMessage) { + await this.appendMessage(sessionId, 'assistant', assistantMessage); + } + + /** + * Elimina el historial de conversación de una sesión. + * @param {string} sessionId - Session ID + * @returns {Promise} + */ + async clearConversation(sessionId) { + if (!Environment.isCacheEnabled(this.envVar, this.native)) { + return; + } + + await redisClient.del(this.getConversationKey(sessionId)); } // Methods implemented for all providers by default @@ -16,7 +206,7 @@ class BaseModel { * @returns {string|undefined} */ getApiKey() { - return process.env[this.apiKeyEnvVar]; + return process.env[`${this.envVar}_API_KEY`]; } /** @@ -24,14 +214,16 @@ class BaseModel { * @returns {string} */ ensureApiKey() { - if (!this.apiKeyEnvVar) { - throw new Error('apiKeyEnvVar is not configured for this provider'); + const envPrefix = this.envVar; + + if (!envPrefix) { + throw new Error('envVar is not configured for this provider'); } const apiKey = this.getApiKey(); if (!apiKey) { - throw new Error(`${this.apiKeyEnvVar} is not configured`); + throw new Error(`${envPrefix}_API_KEY is not configured`); } return apiKey; @@ -110,7 +302,55 @@ class BaseModel { } } - // To be implemented by each provider + /** + * Crea la respuesta del modelo para un mensaje de sesión. + * El flujo común de conversación vive aquí y cada proveedor solo implementa + * la llamada al modelo y el mapeo de respuesta/estado. + * @param {Object} options - Opciones para enviar el mensaje + * @returns {Promise} - Respuesta normalizada del proveedor + */ + async sendMessage(options) { + const { sessionId, message, sessionData } = options; + + if (!sessionId) { + throw Errors.baseModel.missingSessionId(); + } + + const state = new ProviderState(sessionData); + const systemInstruction = state.getSystemInstruction(); + + try { + const conversationMessages = await this.buildConversationForRequest(sessionId, systemInstruction, message); + const requestContext = { + sessionId, + message, + sessionData, + state, + systemInstruction, + conversationMessages, + }; + + const response = await this.buildModelResponse(requestContext); + const responseMessage = this.extractResponseMessage(response, requestContext); + + if (!responseMessage) { + throw Errors.baseModel.noTextContent(this.name); + } + + await this.storeAssistantResponse(sessionId, responseMessage); + + const updatedSessionData = await requestContext.state.buildSessionData(requestContext.sessionId); + + return { + message: responseMessage, + sessionData: updatedSessionData, + }; + } catch (error) { + throw Errors.baseModel.messageSendError(error, this.name); + } + } + + // Methods that each provider must define /** * Crea el cliente del proveedor. Debe ser implementado por cada subclase. @@ -122,16 +362,7 @@ class BaseModel { } /** - * Envía un mensaje a la sesión - * @param {Object} options - Opciones para enviar el mensaje - * @returns {Promise} - Respuesta del modelo - */ - async sendMessage(options) { - throw new Error('Method sendMessage must be implemented by subclasses'); - } - - /** - * Define el threadId para la sesión. Este método debe ser implementado por cada proveedor para determinar cómo manejar el contexto de la conversación. + * Define el threadId para la sesión. Este método es implementado por el proveedor si este nos da el threadId. * @returns {string} El threadId a usar para la sesión, o un string vacío si el proveedor no utiliza threadId. */ async setThreadId() { @@ -147,6 +378,28 @@ class BaseModel { generateEvaluationResponse() { throw new Error('Method generateEvaluationResponse must be implemented by subclasses'); } + + /** + * Ejecuta la llamada al proveedor con el contexto ya preparado. + * Debe implementarse por cada proveedor. + * @param {Object} context - Contexto normalizado de la solicitud + * @returns {Promise} + */ + async buildModelResponse() { + throw new Error('Method buildModelResponse must be implemented by subclasses'); + } + + /** + * Extrae el texto final de la respuesta del proveedor. + * Debe implementarse por cada proveedor. + * @param {Object} response - Respuesta cruda del proveedor + * @param {Object} context - Contexto de la solicitud + * @returns {string} + */ + extractResponseMessage() { + throw new Error('Method extractResponseMessage must be implemented by subclasses'); + } + } module.exports = BaseModel; \ No newline at end of file diff --git a/models/providers/gemini-3.1-flash-lite-preview.js b/models/providers/gemini-3.1-flash-lite-preview.js index 70ef9b7..3976101 100644 --- a/models/providers/gemini-3.1-flash-lite-preview.js +++ b/models/providers/gemini-3.1-flash-lite-preview.js @@ -1,7 +1,6 @@ require('dotenv').config(); const BaseModel = require('./baseModel'); const Errors = require('../../utils/errors'); -const ProviderState = require('../providerState'); const { GoogleGenAI } = require('@google/genai'); /** @@ -12,50 +11,47 @@ class Gemini31FlashLitePreviewProvider extends BaseModel { constructor() { super(); this.name = 'gemini-3.1-flash-lite-preview'; - this.apiKeyEnvVar = 'GEMINI_API_KEY'; + this.envVar = 'GEMINI'; this.model = process.env.GEMINI_MODEL || 'gemini-3.1-flash-lite-preview'; this.evaluationModel = process.env.GEMINI_EVALUATION_MODEL || this.model; } - // Requerido para el baseModel + // Requerido por el baseModel createClient(apiKey) { return new GoogleGenAI({ apiKey }); } - async sendMessage(options) { - const { message, sessionData } = options; - const state = new ProviderState(sessionData); - const systemInstruction = state.getSystemInstruction(); + async buildModelResponse(context) { + const { state, systemInstruction, message } = context; const previousInteractionId = state.get('previousInteractionId') || state.threadId; - try { - const interaction = await this.createInteraction({ - model: this.model, - input: message, - systemInstruction, - previousInteractionId - }); - - const responseMessage = this.extractTextFromInteraction(interaction); - - if (!responseMessage) { - throw Errors.gemini.noTextContent(); - } - - state.update({ - previousInteractionId: interaction.id || previousInteractionId - }); - - return { - message: responseMessage, - sessionData: state.buildSessionData(interaction.id || previousInteractionId), - }; - } catch (error) { - throw Errors.gemini.messageSendError(error); + const interaction = await this.createInteraction({ + model: this.model, + input: message, + systemInstruction, + previousInteractionId, + }); + + context.previousInteractionId = interaction.id || previousInteractionId; + + return interaction; + } + + extractResponseMessage(response) { + if (!response || !Array.isArray(response.outputs)) { + return ''; } + + return response.outputs + .filter(output => output?.type === 'text' && typeof output.text === 'string') + .map(output => output.text.trim()) + .filter(Boolean) + .join('\n\n'); } + + /** * Realiza la llamada al API de Gemini y devuelve la evaluación estructurada. * Invocado por BaseModel.evaluateSolution. @@ -70,7 +66,7 @@ class Gemini31FlashLitePreviewProvider extends BaseModel { responseFormat: this.getEvaluationResponseFormat() }); - const responseText = this.extractTextFromInteraction(interaction); + const responseText = this.extractResponseMessage(interaction); if (!responseText) { throw Errors.gemini.noEvaluationContent(); @@ -105,18 +101,6 @@ class Gemini31FlashLitePreviewProvider extends BaseModel { return fencedMatch ? fencedMatch[1].trim() : trimmedResponse; } - extractTextFromInteraction(interaction) { - if (!interaction || !Array.isArray(interaction.outputs)) { - return ''; - } - - return interaction.outputs - .filter(output => output?.type === 'text' && typeof output.text === 'string') - .map(output => output.text.trim()) - .filter(Boolean) - .join('\n\n'); - } - async createInteraction({ model, input, systemInstruction, previousInteractionId, responseFormat }) { const requestBody = { model, diff --git a/models/providers/ollama.js b/models/providers/ollama.js index 0f32242..547602c 100644 --- a/models/providers/ollama.js +++ b/models/providers/ollama.js @@ -1,23 +1,18 @@ require('dotenv').config(); const BaseModel = require('./baseModel'); const Errors = require('../../utils/errors'); -const ProviderState = require('../providerState'); -const { ConversationStore } = require('../conversationStore'); class OllamaProvider extends BaseModel { constructor() { super(); this.name = 'ollama'; + this.native = false; this.model = process.env.OLLAMA_MODEL || 'llama3.1:8b'; this.evaluationModel = process.env.OLLAMA_EVALUATION_MODEL || this.model; this.baseUrl = (process.env.OLLAMA_BASE_URL || 'http://localhost:11434').replace(/\/+$/, ''); - this.conversationStore = new ConversationStore({ - providerName: 'ollama', - defaultMaxMessages: 60 - }); } - // Requerido por BaseModel + // Requerido por el BaseModel createClient() { return { baseUrl: this.baseUrl, @@ -25,50 +20,29 @@ class OllamaProvider extends BaseModel { }; } - async sendMessage(options) { - const { sessionId, message, sessionData } = options; - - if (!sessionId) { - throw Errors.ollama.missingSessionId(); - } + async buildModelResponse(context) { + const { conversationMessages } = context; - const state = new ProviderState(sessionData); - const systemInstruction = state.getSystemInstruction(); - - try { - const conversationMessages = await this.conversationStore.buildConversationForRequest( - sessionId, - systemInstruction, - message - ); - - const chatResponse = await this.createChatCompletion({ - model: this.model, - messages: conversationMessages, - }); - - const responseMessage = this.extractAssistantMessage(chatResponse); - - if (!responseMessage) { - throw Errors.ollama.noTextContent(); - } + return this.createChatCompletion({ + model: this.model, + messages: conversationMessages, + }); + } - await this.conversationStore.storeAssistantResponse(sessionId, responseMessage); + extractResponseMessage(response) { + if (!response || typeof response !== 'object') { + return ''; + } - state.update({ - conversationKey: this.conversationStore.getConversationKey(sessionId), - model: this.model, - }); + const content = response.message && typeof response.message.content === 'string' + ? response.message.content.trim() + : ''; - return { - message: responseMessage, - sessionData: state.buildSessionData(sessionId), - }; - } catch (error) { - throw Errors.ollama.messageSendError(error); - } + return content; } + + /** * Realiza la llamada al API de Ollama y devuelve la evaluación estructurada. * Invocado por BaseModel.evaluateSolution. @@ -142,18 +116,6 @@ class OllamaProvider extends BaseModel { return responseData; } - extractAssistantMessage(response) { - if (!response || typeof response !== 'object') { - return ''; - } - - const content = response.message && typeof response.message.content === 'string' - ? response.message.content.trim() - : ''; - - return content; - } - getEvaluationResponseFormat() { return { type: 'object', diff --git a/models/providers/openai-responses.js b/models/providers/openai-responses.js index c2b7650..399091e 100644 --- a/models/providers/openai-responses.js +++ b/models/providers/openai-responses.js @@ -4,7 +4,6 @@ const z = require('zod'); const { zodTextFormat } = require('openai/helpers/zod'); const BaseModel = require('./baseModel'); const Errors = require('../../utils/errors'); -const ProviderState = require('../providerState'); const EvaluationSchema = z.object({ score: z.number().min(0).max(10), @@ -18,73 +17,75 @@ class OpenAIResponsesProvider extends BaseModel { constructor() { super(); this.name = 'openai-responses'; - this.apiKeyEnvVar = 'OPENAI_API_KEY'; + this.envVar = 'OPENAI'; this.model = 'gpt-5.4-mini'; this.evaluationModel = process.env.OPENAI_EVALUATION_MODEL || 'gpt-5.4-mini'; } - // Requerido para el baseModel + // Requerido por el BaseModel createClient(apiKey) { return new OpenAI({ apiKey }); } - async sendMessage(options) { - const { message, sessionData } = options; - const state = new ProviderState(sessionData); - const systemInstruction = state.getSystemInstruction(); + async buildModelResponse(context) { + const { sessionData, state, systemInstruction, message } = context; let conversationId = state.get('conversationId') || (state.threadId.startsWith('conv_') ? state.threadId : ''); - try { - if (!conversationId) { - if (sessionData?.threadId) { - console.warn( - 'Sesion legacy de Assistants detectada. Se iniciara una nueva conversacion sin historial previo.' - ); - } - - const conversation = await this.createConversation(); - conversationId = conversation.id; + if (!conversationId) { + if (sessionData?.threadId) { + console.warn( + 'Sesion legacy de Assistants detectada. Se iniciara una nueva conversacion sin historial previo.' + ); } - const response = await this.getClient().responses.create({ - model: this.model, - conversation: conversationId, - instructions: systemInstruction, - input: [ - { - role: 'user', - content: message, - }, - ], - store: true, - }); + const conversation = await this.createConversation(); + conversationId = conversation.id; + } - if (response?.error) { - throw Errors.openAI.responseError(response.error.message); - } + const response = await this.getClient().responses.create({ + model: this.model, + conversation: conversationId, + instructions: systemInstruction, + input: [ + { + role: 'user', + content: message, + }, + ], + store: true, + }); - const responseMessage = this.extractResponseText(response); + if (response?.error) { + throw Errors.openAI.responseError(response.error.message); + } - if (!responseMessage) { - throw Errors.openAI.noTextContent(); - } + context.conversationId = conversationId; + context.lastResponseId = response.id || state.get('lastResponseId'); - state.update({ - conversationId, - systemInstruction, - lastResponseId: response.id || state.get('lastResponseId'), - }); + return response; + } - return { - message: responseMessage, - sessionData: state.buildSessionData(conversationId), - }; - } catch (error) { - throw Errors.openAI.messageSendError(error); + extractResponseMessage(response) { + if (typeof response?.output_text === 'string' && response.output_text.trim()) { + return response.output_text.trim(); + } + + if (!Array.isArray(response?.output)) { + return ''; } + + return response.output + .filter((item) => item?.type === 'message' && Array.isArray(item.content)) + .flatMap((item) => item.content) + .filter((content) => content?.type === 'output_text' && typeof content.text === 'string') + .map((content) => content.text.trim()) + .filter(Boolean) + .join('\n\n'); } + + /** * Realiza la llamada al API de OpenAI y devuelve la evaluación estructurada. * Invocado por BaseModel.evaluateSolution. @@ -134,24 +135,6 @@ class OpenAIResponsesProvider extends BaseModel { return conversation; } - - extractResponseText(response) { - if (typeof response?.output_text === 'string' && response.output_text.trim()) { - return response.output_text.trim(); - } - - if (!Array.isArray(response?.output)) { - return ''; - } - - return response.output - .filter((item) => item?.type === 'message' && Array.isArray(item.content)) - .flatMap((item) => item.content) - .filter((content) => content?.type === 'output_text' && typeof content.text === 'string') - .map((content) => content.text.trim()) - .filter(Boolean) - .join('\n\n'); - } } module.exports = new OpenAIResponsesProvider(); diff --git a/services/cacheService.js b/services/cacheService.js index bf8847e..4074812 100644 --- a/services/cacheService.js +++ b/services/cacheService.js @@ -3,12 +3,27 @@ const { redisClient } = require('../config/redis'); class CacheService { constructor() { this.sessionPrefix = 'session:'; - this.conversationPrefix = 'session:conversation:'; + this.conversationPrefix = 'conversations:'; this.leiaMetaPrefix = 'leia:meta:'; this.modelsPrefix = 'models:'; this.validatedModelsKey = 'validated_models'; } + /** + * Obtiene el sessionId desde una clave de conversación. + * Formato esperado: conversations: + * @param {string} key - Clave de conversación + * @returns {string} sessionId o string vacío si no se puede extraer + */ + extractSessionIdFromConversationKey(key) { + if (!key || !key.startsWith(this.conversationPrefix)) { + return ''; + } + + const parts = key.split(':'); + return parts.length >= 2 ? parts[parts.length - 1] : ''; + } + /** * Parsea el marco temporal a milisegundos * @param {string} timeFrame - Marco temporal (ej: '1h', '2w', '3m', '5d', 'all') @@ -104,7 +119,12 @@ class CacheService { for (const key of keys) { try { if (key.startsWith(this.conversationPrefix)) { - const sessionId = key.replace(this.conversationPrefix, ''); + const sessionId = this.extractSessionIdFromConversationKey(key); + + if (!sessionId) { + continue; + } + const sessionCreatedAt = await redisClient.hGet(`${this.sessionPrefix}${sessionId}`, 'createdAt'); if (sessionCreatedAt) { @@ -143,10 +163,12 @@ class CacheService { * @returns {Array} - Array de claves filtradas */ filterKeysBySession(keys, sessionId) { + const suffix = `:${sessionId}`; + return keys.filter(key => key === `${this.sessionPrefix}${sessionId}` || key === `${this.leiaMetaPrefix}${sessionId}` || - key === `${this.conversationPrefix}${sessionId}` + (key.startsWith(this.conversationPrefix) && key.endsWith(suffix)) ); } @@ -172,8 +194,12 @@ class CacheService { filteredKeys.push(metaKey); } - const conversationKey = `${this.conversationPrefix}${sessionId}`; - if (keys.includes(conversationKey)) { + const conversationSuffix = `:${sessionId}`; + const conversationKeys = keys.filter((keyName) => + keyName.startsWith(this.conversationPrefix) && keyName.endsWith(conversationSuffix) + ); + + for (const conversationKey of conversationKeys) { filteredKeys.push(conversationKey); } } diff --git a/services/sessionService.js b/services/sessionService.js index bd825d2..1f85f28 100644 --- a/services/sessionService.js +++ b/services/sessionService.js @@ -137,6 +137,30 @@ class SessionService { } } + /** + * Clears the cached conversation associated with a session, if the model supports it. + * @param {string} sessionId - Session ID + * @returns {Promise} + */ + async clearConversation(sessionId) { + try { + const sessionData = await this.getSession(sessionId); + + if (!sessionData) { + return; + } + + const model = modelManager.getModel(sessionData.modelName); + + if (model && typeof model.clearConversation === 'function') { + await model.clearConversation(sessionId); + } + } catch (error) { + console.error(`Error clearing conversation cache for session ${sessionId}:`, error); + throw error; + } + } + /** * Stores LEIA metadata associated with the session * @param {string} sessionId - Session ID diff --git a/utils/environment.js b/utils/environment.js new file mode 100644 index 0000000..27e7bd9 --- /dev/null +++ b/utils/environment.js @@ -0,0 +1,53 @@ +/** + * Construye el nombre de una variable de entorno específica del proveedor. + * Ejemplo: OPENAI + API_KEY -> OPENAI_API_KEY. + * @param {string} envVar - Prefijo del proveedor. + * @param {string} suffix - Sufijo de la variable. + * @returns {string} + */ +const getProviderEnvVar = (envVar, suffix) => (envVar ? `${envVar}_${suffix}` : ''); + +/** + * Determina si el cache de conversación está habilitado. + * Prioriza la variable del proveedor y luego la global. + * + * Convención esperada: + * - `${envVar}_CONVERSATION_CACHE_ENABLED` + * - `CONVERSATION_CACHE_ENABLED` + * + * @param {string} envVar - Prefijo del proveedor. + * @param {boolean} native - Indica si el proveedor es nativo. + * @returns {boolean} + */ +function isCacheEnabled(envVar, native = true) { +if (!native) { + return true; +} + +let rawValue; + +if (envVar) { + const providerEnvVar = getProviderEnvVar(envVar, 'CONVERSATION_CACHE_ENABLED'); + rawValue = process.env[providerEnvVar]; +} + +if (rawValue === undefined) { + rawValue = process.env.CONVERSATION_CACHE_ENABLED; +} + +if (rawValue === undefined || rawValue === null || rawValue === '') { + return true; +} + +const normalized = String(rawValue).trim().toLowerCase(); +if (['1', 'true', 'yes', 'on'].includes(normalized)) { + return true; +} +if (['0', 'false', 'no', 'off'].includes(normalized)) { + return false; +} + +return true; +} + +module.exports = { isCacheEnabled }; \ No newline at end of file diff --git a/utils/errors.js b/utils/errors.js index 36fbb49..947e88c 100644 --- a/utils/errors.js +++ b/utils/errors.js @@ -1,6 +1,21 @@ const createError = require('http-errors'); const baseModel = { + missingSessionId: () => + createError(400, 'sessionId es requerido para enviar mensajes al proveedor'), + + noTextContent: (providerName = 'provider') => + createError(500, `${providerName} no devolvio contenido de texto`), + + messageSendError: (originalError, providerName = 'provider') => { + console.error(`Error enviando mensaje a ${providerName}:`, originalError); + if (originalError && (originalError.status || originalError.statusCode)) { + return originalError; + } + + return createError(500, `Error enviando mensaje a ${providerName}`); + }, + missingInstruction: () => createError(500, 'systemInstruction ausente en providerState; posible sesion corrupta'), @@ -14,6 +29,9 @@ const baseModel = { }; const openAI = { + missingSessionId: () => + createError(400, 'sessionId es requerido para usar el proveedor OpenAI'), + noConversationId: () => createError(500, 'OpenAI no devolvio un identificador de conversacion'), @@ -43,6 +61,8 @@ const openAI = { }; const gemini = { + missingSessionId: () => + createError(400, 'sessionId es requerido para usar el proveedor Gemini'), noTextContent: () => createError(500, 'Gemini no devolvio contenido de texto'),