From acd276847fee81f41e400168a3b0ea2b12a14d98 Mon Sep 17 00:00:00 2001 From: Mohamed Ahmed Date: Sat, 18 Apr 2026 20:10:10 +0200 Subject: [PATCH 01/16] feat: new errors --- utils/errors.js | 5 +++++ 1 file changed, 5 insertions(+) diff --git a/utils/errors.js b/utils/errors.js index 36fbb49..970b667 100644 --- a/utils/errors.js +++ b/utils/errors.js @@ -14,6 +14,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 +46,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'), From e3137bb0358b69f7e055e6dd2a6fe5544ad61a0a Mon Sep 17 00:00:00 2001 From: Mohamed Ahmed Date: Sat, 18 Apr 2026 22:07:55 +0200 Subject: [PATCH 02/16] feat: conversation caching for native providers --- .env.example | 9 +++ README.md | 16 +++++ models/conversationStore.js | 72 +++++++++++++++++-- .../gemini-3.1-flash-lite-preview.js | 20 +++++- models/providers/openai-responses.js | 18 ++++- services/cacheService.js | 36 ++++++++-- 6 files changed, 158 insertions(+), 13 deletions(-) diff --git a/.env.example b/.env.example index ab11a13..1d3e397 100644 --- a/.env.example +++ b/.env.example @@ -8,6 +8,15 @@ RUNNER_KEY=R2D2C3PO # Conversation History (Generic) CONVERSATION_HISTORY_MAX_MESSAGES=60 +# Conversation Cache Global switch +# CONVERSATION_CACHE_ENABLED=false + +# Optional provider-specific overrides +# OPENAI_CONVERSATION_CACHE_ENABLED=false +# GEMINI_CONVERSATION_CACHE_ENABLED=false +# OLLAMA_CONVERSATION_CACHE_ENABLED=false +# ... + # Ollama Configuration OLLAMA_BASE_URL=http://localhost:11434 OLLAMA_MODEL=gemma3:4b diff --git a/README.md b/README.md index 96d7c5a..5ae9a46 100644 --- a/README.md +++ b/README.md @@ -10,6 +10,22 @@ 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` +- 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/models/conversationStore.js b/models/conversationStore.js index e9006b3..f00330b 100644 --- a/models/conversationStore.js +++ b/models/conversationStore.js @@ -9,15 +9,43 @@ 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 {string} options.providerName - Provider name for env var lookup and key namespacing * @param {number} options.defaultMaxMessages - Default max messages when not configured (default: 60) */ constructor(options = {}) { - this.keyPrefix = options.prefix || 'conversation:'; - this.providerName = options.providerName || ''; + this.providerName = typeof options.providerName === 'string' ? options.providerName.trim() : ''; + this.basePrefix = 'conversations:'; this.defaultMaxMessages = options.defaultMaxMessages || 60; this.maxMessages = this.parseMaxMessages(); + this.cacheEnabled = this.isCacheEnabled(options.enabled); + } + + /** + * Parses cache enabled flag from explicit option or environment variables. + * Provider-specific env var has priority over global one. + * @private + * @param {boolean|undefined} explicitValue - Explicit option value + * @returns {boolean} Whether cache should be used + */ + isCacheEnabled(explicitValue) { + if (typeof explicitValue === 'boolean') { + return explicitValue; + } + + let rawValue; + + if (this.providerName) { + const providerEnvVar = `${this.providerName.toUpperCase()}_CONVERSATION_CACHE_ENABLED`; + rawValue = process.env[providerEnvVar]; + } + + if (rawValue === undefined) { + rawValue = process.env.CONVERSATION_CACHE_ENABLED; + } + + if (rawValue === undefined || rawValue === null || rawValue === '') { + return true; + } } /** @@ -48,7 +76,9 @@ class ConversationStore { } getConversationKey(sessionId) { - return `${this.keyPrefix}${sessionId}`; + const normalizedSessionId = typeof sessionId === 'string' ? sessionId.trim() : ''; + const providerSegment = this.providerName || 'generic'; + return `${this.basePrefix}${providerSegment}:${normalizedSessionId}`; } /** @@ -81,6 +111,10 @@ class ConversationStore { * @returns {Promise} Array of normalized messages */ async getConversation(sessionId) { + if (!this.cacheEnabled) { + return []; + } + const rawMessages = await redisClient.lRange(this.getConversationKey(sessionId), 0, -1); return rawMessages @@ -109,6 +143,10 @@ class ConversationStore { return; } + if (!this.cacheEnabled) { + return; + } + const key = this.getConversationKey(sessionId); await redisClient.rPush(key, JSON.stringify(message)); await redisClient.lTrim(key, -this.maxMessages, -1); @@ -128,6 +166,10 @@ class ConversationStore { return; } + if (!this.cacheEnabled) { + return; + } + const key = this.getConversationKey(sessionId); const firstRawMessage = await redisClient.lIndex(key, 0); @@ -173,6 +215,22 @@ class ConversationStore { * @returns {Promise} Complete conversation history ready for LLM */ async buildConversationForRequest(sessionId, systemInstruction, userMessage) { + if (!this.cacheEnabled) { + const fallbackConversation = []; + const normalizedSystemMessage = this.normalizeMessage('system', systemInstruction); + const normalizedUserMessage = this.normalizeMessage('user', userMessage); + + if (normalizedSystemMessage) { + fallbackConversation.push(normalizedSystemMessage); + } + + if (normalizedUserMessage) { + fallbackConversation.push(normalizedUserMessage); + } + + return fallbackConversation; + } + await this.ensureSystemMessage(sessionId, systemInstruction); await this.appendMessage(sessionId, 'user', userMessage); return this.getConversation(sessionId); @@ -194,6 +252,10 @@ class ConversationStore { * @returns {Promise} */ async clearConversation(sessionId) { + if (!this.cacheEnabled) { + return; + } + await redisClient.del(this.getConversationKey(sessionId)); } } diff --git a/models/providers/gemini-3.1-flash-lite-preview.js b/models/providers/gemini-3.1-flash-lite-preview.js index 70ef9b7..510fd31 100644 --- a/models/providers/gemini-3.1-flash-lite-preview.js +++ b/models/providers/gemini-3.1-flash-lite-preview.js @@ -2,6 +2,7 @@ require('dotenv').config(); const BaseModel = require('./baseModel'); const Errors = require('../../utils/errors'); const ProviderState = require('../providerState'); +const { ConversationStore } = require('../conversationStore'); const { GoogleGenAI } = require('@google/genai'); /** @@ -15,6 +16,10 @@ class Gemini31FlashLitePreviewProvider extends BaseModel { this.apiKeyEnvVar = 'GEMINI_API_KEY'; this.model = process.env.GEMINI_MODEL || 'gemini-3.1-flash-lite-preview'; this.evaluationModel = process.env.GEMINI_EVALUATION_MODEL || this.model; + this.conversationStore = new ConversationStore({ + providerName: 'gemini', + defaultMaxMessages: 60 + }); } // Requerido para el baseModel @@ -24,12 +29,20 @@ class Gemini31FlashLitePreviewProvider extends BaseModel { } async sendMessage(options) { - const { message, sessionData } = options; + const { sessionId, message, sessionData } = options; + + if (!sessionId) { + throw Errors.gemini.missingSessionId(); + } + const state = new ProviderState(sessionData); const systemInstruction = state.getSystemInstruction(); const previousInteractionId = state.get('previousInteractionId') || state.threadId; try { + await this.conversationStore.ensureSystemMessage(sessionId, systemInstruction); + await this.conversationStore.appendMessage(sessionId, 'user', message); + const interaction = await this.createInteraction({ model: this.model, input: message, @@ -43,8 +56,11 @@ class Gemini31FlashLitePreviewProvider extends BaseModel { throw Errors.gemini.noTextContent(); } + await this.conversationStore.storeAssistantResponse(sessionId, responseMessage); + state.update({ - previousInteractionId: interaction.id || previousInteractionId + previousInteractionId: interaction.id || previousInteractionId, + conversationKey: this.conversationStore.getConversationKey(sessionId) }); return { diff --git a/models/providers/openai-responses.js b/models/providers/openai-responses.js index c2b7650..f729a62 100644 --- a/models/providers/openai-responses.js +++ b/models/providers/openai-responses.js @@ -5,6 +5,7 @@ const { zodTextFormat } = require('openai/helpers/zod'); const BaseModel = require('./baseModel'); const Errors = require('../../utils/errors'); const ProviderState = require('../providerState'); +const { ConversationStore } = require('../conversationStore'); const EvaluationSchema = z.object({ score: z.number().min(0).max(10), @@ -21,6 +22,10 @@ class OpenAIResponsesProvider extends BaseModel { this.apiKeyEnvVar = 'OPENAI_API_KEY'; this.model = 'gpt-5.4-mini'; this.evaluationModel = process.env.OPENAI_EVALUATION_MODEL || 'gpt-5.4-mini'; + this.conversationStore = new ConversationStore({ + providerName: 'openai', + defaultMaxMessages: 60, + }); } // Requerido para el baseModel @@ -30,7 +35,12 @@ class OpenAIResponsesProvider extends BaseModel { } async sendMessage(options) { - const { message, sessionData } = options; + const { sessionId, message, sessionData } = options; + + if (!sessionId) { + throw Errors.openAI.missingSessionId(); + } + const state = new ProviderState(sessionData); const systemInstruction = state.getSystemInstruction(); let conversationId = state.get('conversationId') || (state.threadId.startsWith('conv_') ? state.threadId : ''); @@ -47,6 +57,9 @@ class OpenAIResponsesProvider extends BaseModel { conversationId = conversation.id; } + await this.conversationStore.ensureSystemMessage(sessionId, systemInstruction); + await this.conversationStore.appendMessage(sessionId, 'user', message); + const response = await this.getClient().responses.create({ model: this.model, conversation: conversationId, @@ -70,8 +83,11 @@ class OpenAIResponsesProvider extends BaseModel { throw Errors.openAI.noTextContent(); } + await this.conversationStore.storeAssistantResponse(sessionId, responseMessage); + state.update({ conversationId, + conversationKey: this.conversationStore.getConversationKey(sessionId), systemInstruction, lastResponseId: response.id || state.get('lastResponseId'), }); diff --git a/services/cacheService.js b/services/cacheService.js index bf8847e..d894f5e 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 >= 3 ? 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); } } From 450bdb9b652e16dd7ca3e13c2760460f8d4022bd Mon Sep 17 00:00:00 2001 From: Mohamed Ahmed Date: Thu, 23 Apr 2026 13:07:13 +0200 Subject: [PATCH 03/16] fix: env.example --- .env.example | 8 ++++---- 1 file changed, 4 insertions(+), 4 deletions(-) diff --git a/.env.example b/.env.example index 1d3e397..82cec07 100644 --- a/.env.example +++ b/.env.example @@ -9,12 +9,12 @@ RUNNER_KEY=R2D2C3PO CONVERSATION_HISTORY_MAX_MESSAGES=60 # Conversation Cache Global switch -# CONVERSATION_CACHE_ENABLED=false +CONVERSATION_CACHE_ENABLED=true # Optional provider-specific overrides -# OPENAI_CONVERSATION_CACHE_ENABLED=false -# GEMINI_CONVERSATION_CACHE_ENABLED=false -# OLLAMA_CONVERSATION_CACHE_ENABLED=false +OPENAI_CONVERSATION_CACHE_ENABLED=true +GEMINI_CONVERSATION_CACHE_ENABLED=true +OLLAMA_CONVERSATION_CACHE_ENABLED=true # ... # Ollama Configuration From b2930697d183e20104db1cf05fdf165350a509ac Mon Sep 17 00:00:00 2001 From: Mohamed Ahmed Date: Thu, 23 Apr 2026 15:34:51 +0200 Subject: [PATCH 04/16] fix: conversaciones en cache no segregadas --- models/conversationStore.js | 3 +-- services/cacheService.js | 4 ++-- 2 files changed, 3 insertions(+), 4 deletions(-) diff --git a/models/conversationStore.js b/models/conversationStore.js index f00330b..60f2bbf 100644 --- a/models/conversationStore.js +++ b/models/conversationStore.js @@ -77,8 +77,7 @@ class ConversationStore { getConversationKey(sessionId) { const normalizedSessionId = typeof sessionId === 'string' ? sessionId.trim() : ''; - const providerSegment = this.providerName || 'generic'; - return `${this.basePrefix}${providerSegment}:${normalizedSessionId}`; + return `${this.basePrefix}${normalizedSessionId}`; } /** diff --git a/services/cacheService.js b/services/cacheService.js index d894f5e..4074812 100644 --- a/services/cacheService.js +++ b/services/cacheService.js @@ -11,7 +11,7 @@ class CacheService { /** * Obtiene el sessionId desde una clave de conversación. - * Formato esperado: conversations:: + * Formato esperado: conversations: * @param {string} key - Clave de conversación * @returns {string} sessionId o string vacío si no se puede extraer */ @@ -21,7 +21,7 @@ class CacheService { } const parts = key.split(':'); - return parts.length >= 3 ? parts[parts.length - 1] : ''; + return parts.length >= 2 ? parts[parts.length - 1] : ''; } /** From 39b772328c06f6956b5e3ca88380d27b1c674dca Mon Sep 17 00:00:00 2001 From: Mohamed Ahmed Date: Thu, 23 Apr 2026 15:48:07 +0200 Subject: [PATCH 05/16] feat: conversations now are cleared when solution is sended Co-authored-by: Copilot --- controllers/evaluationController.js | 12 ++++++++++++ services/sessionService.js | 25 +++++++++++++++++++++++++ 2 files changed, 37 insertions(+) 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/services/sessionService.js b/services/sessionService.js index bd825d2..ce8aaab 100644 --- a/services/sessionService.js +++ b/services/sessionService.js @@ -137,6 +137,31 @@ 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); + const conversationStore = model?.conversationStore; + + if (conversationStore && typeof conversationStore.clearConversation === 'function') { + await conversationStore.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 From a0d5aa83efdacc462179f2e343a98b454efaa385 Mon Sep 17 00:00:00 2001 From: Mohamed Ahmed Date: Thu, 23 Apr 2026 16:39:55 +0200 Subject: [PATCH 06/16] feat: add TTL for conversation cache Co-authored-by: Copilot --- .env.example | 1 + README.md | 1 + models/conversationStore.js | 17 +++++++++++++++++ 3 files changed, 19 insertions(+) diff --git a/.env.example b/.env.example index 82cec07..ecc173a 100644 --- a/.env.example +++ b/.env.example @@ -7,6 +7,7 @@ RUNNER_KEY=R2D2C3PO # Conversation History (Generic) CONVERSATION_HISTORY_MAX_MESSAGES=60 +CONVERSATION_HISTORY_TTL=2629800 # Conversation Cache Global switch CONVERSATION_CACHE_ENABLED=true diff --git a/README.md b/README.md index 5ae9a46..eae5420 100644 --- a/README.md +++ b/README.md @@ -15,6 +15,7 @@ API for interacting with LEIA instances. 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` diff --git a/models/conversationStore.js b/models/conversationStore.js index 60f2bbf..a3ff294 100644 --- a/models/conversationStore.js +++ b/models/conversationStore.js @@ -17,6 +17,7 @@ class ConversationStore { this.basePrefix = 'conversations:'; this.defaultMaxMessages = options.defaultMaxMessages || 60; this.maxMessages = this.parseMaxMessages(); + this.ttlSeconds = process.env.CONVERSATION_HISTORY_TTL || 2629800; this.cacheEnabled = this.isCacheEnabled(options.enabled); } @@ -75,6 +76,16 @@ class ConversationStore { return parsed; } + /** + * Refreshes the TTL for a conversation key. + * @private + * @param {string} key - Redis conversation key + * @returns {Promise} + */ + async refreshConversationTtl(key) { + await redisClient.expire(key, this.ttlSeconds); + } + getConversationKey(sessionId) { const normalizedSessionId = typeof sessionId === 'string' ? sessionId.trim() : ''; return `${this.basePrefix}${normalizedSessionId}`; @@ -149,6 +160,7 @@ class ConversationStore { const key = this.getConversationKey(sessionId); await redisClient.rPush(key, JSON.stringify(message)); await redisClient.lTrim(key, -this.maxMessages, -1); + await this.refreshConversationTtl(key); } /** @@ -174,6 +186,7 @@ class ConversationStore { if (!firstRawMessage) { await redisClient.rPush(key, JSON.stringify(normalizedSystemMessage)); + await this.refreshConversationTtl(key); return; } @@ -191,18 +204,22 @@ class ConversationStore { if (!normalizedFirstMessage) { await redisClient.lSet(key, 0, JSON.stringify(normalizedSystemMessage)); + await this.refreshConversationTtl(key); return; } if (normalizedFirstMessage.role !== 'system') { await redisClient.lPush(key, JSON.stringify(normalizedSystemMessage)); await redisClient.lTrim(key, -this.maxMessages, -1); + await this.refreshConversationTtl(key); return; } if (normalizedFirstMessage.content !== normalizedSystemMessage.content) { await redisClient.lSet(key, 0, JSON.stringify(normalizedSystemMessage)); } + + await this.refreshConversationTtl(key); } /** From 6bf7840f198365e43ae9b9d5487299322589ac99 Mon Sep 17 00:00:00 2001 From: Mohamed Ahmed Date: Thu, 23 Apr 2026 17:39:45 +0200 Subject: [PATCH 07/16] refactor: migrate conversation management to BaseModel for all providers Co-authored-by: Copilot --- models/conversationStore.js | 279 ------------------ models/providers/baseModel.js | 279 ++++++++++++++++++ .../gemini-3.1-flash-lite-preview.js | 13 +- models/providers/ollama.js | 11 +- models/providers/openai-responses.js | 13 +- services/sessionService.js | 5 +- 6 files changed, 292 insertions(+), 308 deletions(-) delete mode 100644 models/conversationStore.js diff --git a/models/conversationStore.js b/models/conversationStore.js deleted file mode 100644 index a3ff294..0000000 --- a/models/conversationStore.js +++ /dev/null @@ -1,279 +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.providerName - Provider name for env var lookup and key namespacing - * @param {number} options.defaultMaxMessages - Default max messages when not configured (default: 60) - */ - constructor(options = {}) { - this.providerName = typeof options.providerName === 'string' ? options.providerName.trim() : ''; - this.basePrefix = 'conversations:'; - this.defaultMaxMessages = options.defaultMaxMessages || 60; - this.maxMessages = this.parseMaxMessages(); - this.ttlSeconds = process.env.CONVERSATION_HISTORY_TTL || 2629800; - this.cacheEnabled = this.isCacheEnabled(options.enabled); - } - - /** - * Parses cache enabled flag from explicit option or environment variables. - * Provider-specific env var has priority over global one. - * @private - * @param {boolean|undefined} explicitValue - Explicit option value - * @returns {boolean} Whether cache should be used - */ - isCacheEnabled(explicitValue) { - if (typeof explicitValue === 'boolean') { - return explicitValue; - } - - let rawValue; - - if (this.providerName) { - const providerEnvVar = `${this.providerName.toUpperCase()}_CONVERSATION_CACHE_ENABLED`; - rawValue = process.env[providerEnvVar]; - } - - if (rawValue === undefined) { - rawValue = process.env.CONVERSATION_CACHE_ENABLED; - } - - if (rawValue === undefined || rawValue === null || rawValue === '') { - return true; - } - } - - /** - * 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; - } - - /** - * Refreshes the TTL for a conversation key. - * @private - * @param {string} key - Redis conversation key - * @returns {Promise} - */ - async refreshConversationTtl(key) { - await redisClient.expire(key, this.ttlSeconds); - } - - getConversationKey(sessionId) { - const normalizedSessionId = typeof sessionId === 'string' ? sessionId.trim() : ''; - return `${this.basePrefix}${normalizedSessionId}`; - } - - /** - * 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) { - if (!this.cacheEnabled) { - 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.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; - } - - if (!this.cacheEnabled) { - return; - } - - const key = this.getConversationKey(sessionId); - await redisClient.rPush(key, JSON.stringify(message)); - await redisClient.lTrim(key, -this.maxMessages, -1); - await this.refreshConversationTtl(key); - } - - /** - * 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; - } - - if (!this.cacheEnabled) { - return; - } - - const key = this.getConversationKey(sessionId); - const firstRawMessage = await redisClient.lIndex(key, 0); - - if (!firstRawMessage) { - await redisClient.rPush(key, JSON.stringify(normalizedSystemMessage)); - await this.refreshConversationTtl(key); - 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)); - await this.refreshConversationTtl(key); - return; - } - - if (normalizedFirstMessage.role !== 'system') { - await redisClient.lPush(key, JSON.stringify(normalizedSystemMessage)); - await redisClient.lTrim(key, -this.maxMessages, -1); - await this.refreshConversationTtl(key); - return; - } - - if (normalizedFirstMessage.content !== normalizedSystemMessage.content) { - await redisClient.lSet(key, 0, JSON.stringify(normalizedSystemMessage)); - } - - await this.refreshConversationTtl(key); - } - - /** - * 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) { - if (!this.cacheEnabled) { - const fallbackConversation = []; - const normalizedSystemMessage = this.normalizeMessage('system', systemInstruction); - const normalizedUserMessage = this.normalizeMessage('user', userMessage); - - if (normalizedSystemMessage) { - fallbackConversation.push(normalizedSystemMessage); - } - - if (normalizedUserMessage) { - fallbackConversation.push(normalizedUserMessage); - } - - return fallbackConversation; - } - - 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) { - if (!this.cacheEnabled) { - return; - } - - 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..23d7c5e 100644 --- a/models/providers/baseModel.js +++ b/models/providers/baseModel.js @@ -1,12 +1,291 @@ const Errors = require('../../utils/errors'); const Prompts = require('../../utils/prompts'); const ProviderState = require('../providerState'); +const { redisClient } = require('../../config/redis'); class BaseModel { constructor() { this.name = 'base'; this.apiKeyEnvVar = ''; this._client = null; + this.conversationPrefix = 'conversations:'; + this.defaultConversationMaxMessages = 60; + 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 el prefijo de variables de entorno para configuración de conversación. + * @returns {string} + */ + getConversationEnvPrefix() { + const apiKeyEnvVar = typeof this.apiKeyEnvVar === 'string' ? this.apiKeyEnvVar.trim() : ''; + + if (apiKeyEnvVar.endsWith('_API_KEY')) { + return apiKeyEnvVar.slice(0, -'_API_KEY'.length); + } + + const providerName = typeof this.name === 'string' ? this.name.trim() : ''; + return providerName.toUpperCase().replace(/[^A-Z0-9]+/g, '_'); + } + + /** + * Determina si el cache de conversación está habilitado. + * Prioriza variable por proveedor y luego la global. + * @returns {boolean} + */ + isConversationCacheEnabled() { + const envPrefix = this.getConversationEnvPrefix(); + let rawValue; + + if (envPrefix) { + const providerEnvVar = `${envPrefix}_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; + } + + /** + * Obtiene el límite de historial de conversación para el proveedor actual. + * @returns {number} + */ + getConversationMaxMessages() { + const envPrefix = this.getConversationEnvPrefix(); + let rawValue; + + if (envPrefix) { + const providerEnvVar = `${envPrefix}_HISTORY_MAX_MESSAGES`; + rawValue = process.env[providerEnvVar]; + } + + if (!rawValue) { + rawValue = process.env.CONVERSATION_HISTORY_MAX_MESSAGES; + } + + const parsed = Number.parseInt(rawValue || this.defaultConversationMaxMessages, 10); + if (!Number.isInteger(parsed) || parsed <= 0) { + return this.defaultConversationMaxMessages; + } + + return parsed; + } + + /** + * Obtiene TTL de historial de conversación en segundos. + * @returns {number} + */ + getConversationTtlSeconds() { + const parsed = Number.parseInt(process.env.CONVERSATION_HISTORY_TTL || this.defaultConversationTtlSeconds, 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 (!this.isConversationCacheEnabled()) { + 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 || !this.isConversationCacheEnabled()) { + return; + } + + const key = this.getConversationKey(sessionId); + const maxMessages = this.getConversationMaxMessages(); + await redisClient.rPush(key, JSON.stringify(message)); + await redisClient.lTrim(key, -maxMessages, -1); + 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 || !this.isConversationCacheEnabled()) { + return; + } + + const key = this.getConversationKey(sessionId); + const maxMessages = this.getConversationMaxMessages(); + 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.lTrim(key, -maxMessages, -1); + 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) { + if (!this.isConversationCacheEnabled()) { + const fallbackConversation = []; + const normalizedSystemMessage = this.normalizeConversationMessage('system', systemInstruction); + const normalizedUserMessage = this.normalizeConversationMessage('user', userMessage); + + if (normalizedSystemMessage) { + fallbackConversation.push(normalizedSystemMessage); + } + + if (normalizedUserMessage) { + fallbackConversation.push(normalizedUserMessage); + } + + return fallbackConversation; + } + + 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 (!this.isConversationCacheEnabled()) { + return; + } + + await redisClient.del(this.getConversationKey(sessionId)); } // Methods implemented for all providers by default diff --git a/models/providers/gemini-3.1-flash-lite-preview.js b/models/providers/gemini-3.1-flash-lite-preview.js index 510fd31..53870a6 100644 --- a/models/providers/gemini-3.1-flash-lite-preview.js +++ b/models/providers/gemini-3.1-flash-lite-preview.js @@ -2,7 +2,6 @@ require('dotenv').config(); const BaseModel = require('./baseModel'); const Errors = require('../../utils/errors'); const ProviderState = require('../providerState'); -const { ConversationStore } = require('../conversationStore'); const { GoogleGenAI } = require('@google/genai'); /** @@ -16,10 +15,6 @@ class Gemini31FlashLitePreviewProvider extends BaseModel { this.apiKeyEnvVar = 'GEMINI_API_KEY'; this.model = process.env.GEMINI_MODEL || 'gemini-3.1-flash-lite-preview'; this.evaluationModel = process.env.GEMINI_EVALUATION_MODEL || this.model; - this.conversationStore = new ConversationStore({ - providerName: 'gemini', - defaultMaxMessages: 60 - }); } // Requerido para el baseModel @@ -40,8 +35,8 @@ class Gemini31FlashLitePreviewProvider extends BaseModel { const previousInteractionId = state.get('previousInteractionId') || state.threadId; try { - await this.conversationStore.ensureSystemMessage(sessionId, systemInstruction); - await this.conversationStore.appendMessage(sessionId, 'user', message); + await this.ensureSystemMessage(sessionId, systemInstruction); + await this.appendMessage(sessionId, 'user', message); const interaction = await this.createInteraction({ model: this.model, @@ -56,11 +51,11 @@ class Gemini31FlashLitePreviewProvider extends BaseModel { throw Errors.gemini.noTextContent(); } - await this.conversationStore.storeAssistantResponse(sessionId, responseMessage); + await this.storeAssistantResponse(sessionId, responseMessage); state.update({ previousInteractionId: interaction.id || previousInteractionId, - conversationKey: this.conversationStore.getConversationKey(sessionId) + conversationKey: this.getConversationKey(sessionId) }); return { diff --git a/models/providers/ollama.js b/models/providers/ollama.js index 0f32242..181a636 100644 --- a/models/providers/ollama.js +++ b/models/providers/ollama.js @@ -2,7 +2,6 @@ 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() { @@ -11,10 +10,6 @@ class OllamaProvider extends BaseModel { 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 @@ -36,7 +31,7 @@ class OllamaProvider extends BaseModel { const systemInstruction = state.getSystemInstruction(); try { - const conversationMessages = await this.conversationStore.buildConversationForRequest( + const conversationMessages = await this.buildConversationForRequest( sessionId, systemInstruction, message @@ -53,10 +48,10 @@ class OllamaProvider extends BaseModel { throw Errors.ollama.noTextContent(); } - await this.conversationStore.storeAssistantResponse(sessionId, responseMessage); + await this.storeAssistantResponse(sessionId, responseMessage); state.update({ - conversationKey: this.conversationStore.getConversationKey(sessionId), + conversationKey: this.getConversationKey(sessionId), model: this.model, }); diff --git a/models/providers/openai-responses.js b/models/providers/openai-responses.js index f729a62..6d1992a 100644 --- a/models/providers/openai-responses.js +++ b/models/providers/openai-responses.js @@ -5,7 +5,6 @@ const { zodTextFormat } = require('openai/helpers/zod'); const BaseModel = require('./baseModel'); const Errors = require('../../utils/errors'); const ProviderState = require('../providerState'); -const { ConversationStore } = require('../conversationStore'); const EvaluationSchema = z.object({ score: z.number().min(0).max(10), @@ -22,10 +21,6 @@ class OpenAIResponsesProvider extends BaseModel { this.apiKeyEnvVar = 'OPENAI_API_KEY'; this.model = 'gpt-5.4-mini'; this.evaluationModel = process.env.OPENAI_EVALUATION_MODEL || 'gpt-5.4-mini'; - this.conversationStore = new ConversationStore({ - providerName: 'openai', - defaultMaxMessages: 60, - }); } // Requerido para el baseModel @@ -57,8 +52,8 @@ class OpenAIResponsesProvider extends BaseModel { conversationId = conversation.id; } - await this.conversationStore.ensureSystemMessage(sessionId, systemInstruction); - await this.conversationStore.appendMessage(sessionId, 'user', message); + await this.ensureSystemMessage(sessionId, systemInstruction); + await this.appendMessage(sessionId, 'user', message); const response = await this.getClient().responses.create({ model: this.model, @@ -83,11 +78,11 @@ class OpenAIResponsesProvider extends BaseModel { throw Errors.openAI.noTextContent(); } - await this.conversationStore.storeAssistantResponse(sessionId, responseMessage); + await this.storeAssistantResponse(sessionId, responseMessage); state.update({ conversationId, - conversationKey: this.conversationStore.getConversationKey(sessionId), + conversationKey: this.getConversationKey(sessionId), systemInstruction, lastResponseId: response.id || state.get('lastResponseId'), }); diff --git a/services/sessionService.js b/services/sessionService.js index ce8aaab..1f85f28 100644 --- a/services/sessionService.js +++ b/services/sessionService.js @@ -151,10 +151,9 @@ class SessionService { } const model = modelManager.getModel(sessionData.modelName); - const conversationStore = model?.conversationStore; - if (conversationStore && typeof conversationStore.clearConversation === 'function') { - await conversationStore.clearConversation(sessionId); + if (model && typeof model.clearConversation === 'function') { + await model.clearConversation(sessionId); } } catch (error) { console.error(`Error clearing conversation cache for session ${sessionId}:`, error); From 5ea3965c5cd4ebc1a194e0c9721f2d111d925e31 Mon Sep 17 00:00:00 2001 From: Mohamed Ahmed Date: Sun, 26 Apr 2026 00:32:15 +0200 Subject: [PATCH 08/16] refactor: replace apiKeyEnvVar with envVar in BaseModel and providers; add environment utility functions Co-authored-by: Copilot --- models/providers/baseModel.js | 101 +++--------------- .../gemini-3.1-flash-lite-preview.js | 2 +- models/providers/openai-responses.js | 2 +- utils/environment.js | 79 ++++++++++++++ 4 files changed, 96 insertions(+), 88 deletions(-) create mode 100644 utils/environment.js diff --git a/models/providers/baseModel.js b/models/providers/baseModel.js index 23d7c5e..6741668 100644 --- a/models/providers/baseModel.js +++ b/models/providers/baseModel.js @@ -1,19 +1,19 @@ 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.envVar = null; this._client = null; this.conversationPrefix = 'conversations:'; this.defaultConversationMaxMessages = 60; this.defaultConversationTtlSeconds = 2629800; } - // conversationStore methods implemented for all providers by default /** @@ -26,79 +26,6 @@ class BaseModel { return `${this.conversationPrefix}${normalizedSessionId}`; } - /** - * Obtiene el prefijo de variables de entorno para configuración de conversación. - * @returns {string} - */ - getConversationEnvPrefix() { - const apiKeyEnvVar = typeof this.apiKeyEnvVar === 'string' ? this.apiKeyEnvVar.trim() : ''; - - if (apiKeyEnvVar.endsWith('_API_KEY')) { - return apiKeyEnvVar.slice(0, -'_API_KEY'.length); - } - - const providerName = typeof this.name === 'string' ? this.name.trim() : ''; - return providerName.toUpperCase().replace(/[^A-Z0-9]+/g, '_'); - } - - /** - * Determina si el cache de conversación está habilitado. - * Prioriza variable por proveedor y luego la global. - * @returns {boolean} - */ - isConversationCacheEnabled() { - const envPrefix = this.getConversationEnvPrefix(); - let rawValue; - - if (envPrefix) { - const providerEnvVar = `${envPrefix}_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; - } - - /** - * Obtiene el límite de historial de conversación para el proveedor actual. - * @returns {number} - */ - getConversationMaxMessages() { - const envPrefix = this.getConversationEnvPrefix(); - let rawValue; - - if (envPrefix) { - const providerEnvVar = `${envPrefix}_HISTORY_MAX_MESSAGES`; - rawValue = process.env[providerEnvVar]; - } - - if (!rawValue) { - rawValue = process.env.CONVERSATION_HISTORY_MAX_MESSAGES; - } - - const parsed = Number.parseInt(rawValue || this.defaultConversationMaxMessages, 10); - if (!Number.isInteger(parsed) || parsed <= 0) { - return this.defaultConversationMaxMessages; - } - - return parsed; - } - /** * Obtiene TTL de historial de conversación en segundos. * @returns {number} @@ -142,7 +69,7 @@ class BaseModel { * @returns {Promise} */ async getConversation(sessionId) { - if (!this.isConversationCacheEnabled()) { + if (!Environment.isCacheEnabled(this.envVar)) { return []; } @@ -170,12 +97,12 @@ class BaseModel { async appendMessage(sessionId, role, content) { const message = this.normalizeConversationMessage(role, content); - if (!message || !this.isConversationCacheEnabled()) { + if (!message || !Environment.isCacheEnabled(this.envVar)) { return; } const key = this.getConversationKey(sessionId); - const maxMessages = this.getConversationMaxMessages(); + const maxMessages = Environment.getConversationMaxMessages(this.envVar); await redisClient.rPush(key, JSON.stringify(message)); await redisClient.lTrim(key, -maxMessages, -1); await redisClient.expire(key, this.getConversationTtlSeconds()); @@ -190,12 +117,12 @@ class BaseModel { async ensureSystemMessage(sessionId, systemInstruction) { const normalizedSystemMessage = this.normalizeConversationMessage('system', systemInstruction); - if (!normalizedSystemMessage || !this.isConversationCacheEnabled()) { + if (!normalizedSystemMessage || !Environment.isCacheEnabled(this.envVar)) { return; } const key = this.getConversationKey(sessionId); - const maxMessages = this.getConversationMaxMessages(); + const maxMessages = Environment.getConversationMaxMessages(this.envVar); const firstRawMessage = await redisClient.lIndex(key, 0); if (!firstRawMessage) { @@ -244,7 +171,7 @@ class BaseModel { * @returns {Promise} */ async buildConversationForRequest(sessionId, systemInstruction, userMessage) { - if (!this.isConversationCacheEnabled()) { + if (!Environment.isCacheEnabled(this.envVar)) { const fallbackConversation = []; const normalizedSystemMessage = this.normalizeConversationMessage('system', systemInstruction); const normalizedUserMessage = this.normalizeConversationMessage('user', userMessage); @@ -281,7 +208,7 @@ class BaseModel { * @returns {Promise} */ async clearConversation(sessionId) { - if (!this.isConversationCacheEnabled()) { + if (!Environment.isCacheEnabled(this.envVar)) { return; } @@ -295,7 +222,7 @@ class BaseModel { * @returns {string|undefined} */ getApiKey() { - return process.env[this.apiKeyEnvVar]; + return process.env[`${this.envVar}_API_KEY`]; } /** @@ -303,14 +230,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; diff --git a/models/providers/gemini-3.1-flash-lite-preview.js b/models/providers/gemini-3.1-flash-lite-preview.js index 53870a6..b45e2cd 100644 --- a/models/providers/gemini-3.1-flash-lite-preview.js +++ b/models/providers/gemini-3.1-flash-lite-preview.js @@ -12,7 +12,7 @@ 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; } diff --git a/models/providers/openai-responses.js b/models/providers/openai-responses.js index 6d1992a..63d875c 100644 --- a/models/providers/openai-responses.js +++ b/models/providers/openai-responses.js @@ -18,7 +18,7 @@ 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'; } diff --git a/utils/environment.js b/utils/environment.js new file mode 100644 index 0000000..d7a1bad --- /dev/null +++ b/utils/environment.js @@ -0,0 +1,79 @@ +/** + * 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. + * @returns {boolean} + */ +function isCacheEnabled(envVar) { +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; +} + +/** + * Obtiene el límite de historial de conversación para el proveedor actual. + * Prioriza la variable del proveedor y luego la global. + * + * Convención esperada: + * - `${envVar}_HISTORY_MAX_MESSAGES` + * - `CONVERSATION_HISTORY_MAX_MESSAGES` + * + * @param {string} envVar - Prefijo del proveedor. + * @returns {number} + */ +function getConversationMaxMessages(envVar) { +let rawValue; + +if (envVar) { + const providerEnvVar = getProviderEnvVar(envVar, 'HISTORY_MAX_MESSAGES'); + rawValue = process.env[providerEnvVar]; +} + +if (!rawValue) { + rawValue = process.env.CONVERSATION_HISTORY_MAX_MESSAGES; +} + +const parsed = Number.parseInt(rawValue || 60, 10); +if (!Number.isInteger(parsed) || parsed <= 0) { + return 60; +} + +return parsed; +} + +module.exports = { isCacheEnabled, getConversationMaxMessages }; \ No newline at end of file From 5be1223050906811d08872d34020cf30caf86b72 Mon Sep 17 00:00:00 2001 From: Mohamed Ahmed Date: Sun, 26 Apr 2026 13:44:31 +0200 Subject: [PATCH 09/16] feat: enhance BaseModel with improved message handling and error management; add new methods for response processing Co-authored-by: Copilot --- models/providers/baseModel.js | 101 ++++++++++++++++-- .../gemini-3.1-flash-lite-preview.js | 58 ++++------ models/providers/ollama.js | 57 ++++------ models/providers/openai-responses.js | 99 ++++++++--------- utils/errors.js | 15 +++ 5 files changed, 188 insertions(+), 142 deletions(-) diff --git a/models/providers/baseModel.js b/models/providers/baseModel.js index 6741668..3708775 100644 --- a/models/providers/baseModel.js +++ b/models/providers/baseModel.js @@ -31,7 +31,11 @@ class BaseModel { * @returns {number} */ getConversationTtlSeconds() { - const parsed = Number.parseInt(process.env.CONVERSATION_HISTORY_TTL || this.defaultConversationTtlSeconds, 10); + 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; } @@ -318,7 +322,59 @@ 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 this.buildSessionDataAfterMessage( + requestContext, + response, + responseMessage + ); + + 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. @@ -329,15 +385,6 @@ class BaseModel { throw new Error('Method createClient must be implemented by subclasses'); } - /** - * 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. * @returns {string} El threadId a usar para la sesión, o un string vacío si el proveedor no utiliza threadId. @@ -355,6 +402,38 @@ 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'); + } + + /** + * Construye el sessionData que se devolverá al servicio de sesión. + * @param {Object} context - Contexto de la solicitud + * @param {Object} response - Respuesta cruda del proveedor + * @param {string} responseMessage - Mensaje final del asistente + * @returns {Object} + */ + async buildSessionDataAfterMessage(context, response, responseMessage) { + return context.state.buildSessionData(context.sessionId); + } } 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 b45e2cd..d38a72e 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'); /** @@ -17,54 +16,41 @@ class Gemini31FlashLitePreviewProvider extends BaseModel { 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 { sessionId, message, sessionData } = options; - - if (!sessionId) { - throw Errors.gemini.missingSessionId(); - } - - 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 { - await this.ensureSystemMessage(sessionId, systemInstruction); - await this.appendMessage(sessionId, 'user', message); + const interaction = await this.createInteraction({ + model: this.model, + input: message, + systemInstruction, + previousInteractionId, + }); - const interaction = await this.createInteraction({ - model: this.model, - input: message, - systemInstruction, - previousInteractionId - }); + context.previousInteractionId = interaction.id || previousInteractionId; - const responseMessage = this.extractTextFromInteraction(interaction); + return interaction; + } - if (!responseMessage) { - throw Errors.gemini.noTextContent(); - } + extractResponseMessage(response) { + return this.extractTextFromInteraction(response); + } - await this.storeAssistantResponse(sessionId, responseMessage); + async buildSessionDataAfterMessage(context) { + const { state, sessionId, previousInteractionId } = context; - state.update({ - previousInteractionId: interaction.id || previousInteractionId, - conversationKey: this.getConversationKey(sessionId) - }); + state.update({ + previousInteractionId, + conversationKey: this.getConversationKey(sessionId), + }); - return { - message: responseMessage, - sessionData: state.buildSessionData(interaction.id || previousInteractionId), - }; - } catch (error) { - throw Errors.gemini.messageSendError(error); - } + return state.buildSessionData(previousInteractionId); } /** diff --git a/models/providers/ollama.js b/models/providers/ollama.js index 181a636..cdaf8ac 100644 --- a/models/providers/ollama.js +++ b/models/providers/ollama.js @@ -1,7 +1,6 @@ require('dotenv').config(); const BaseModel = require('./baseModel'); const Errors = require('../../utils/errors'); -const ProviderState = require('../providerState'); class OllamaProvider extends BaseModel { constructor() { @@ -12,7 +11,7 @@ class OllamaProvider extends BaseModel { this.baseUrl = (process.env.OLLAMA_BASE_URL || 'http://localhost:11434').replace(/\/+$/, ''); } - // Requerido por BaseModel + // Requerido por el BaseModel createClient() { return { baseUrl: this.baseUrl, @@ -20,48 +19,28 @@ class OllamaProvider extends BaseModel { }; } - async sendMessage(options) { - const { sessionId, message, sessionData } = options; + async buildModelResponse(context) { + const { conversationMessages } = context; - if (!sessionId) { - throw Errors.ollama.missingSessionId(); - } - - const state = new ProviderState(sessionData); - const systemInstruction = state.getSystemInstruction(); - - try { - const conversationMessages = await this.buildConversationForRequest( - sessionId, - systemInstruction, - message - ); - - const chatResponse = await this.createChatCompletion({ - model: this.model, - messages: conversationMessages, - }); - - const responseMessage = this.extractAssistantMessage(chatResponse); + return this.createChatCompletion({ + model: this.model, + messages: conversationMessages, + }); + } - if (!responseMessage) { - throw Errors.ollama.noTextContent(); - } + extractResponseMessage(response) { + return this.extractAssistantMessage(response); + } - await this.storeAssistantResponse(sessionId, responseMessage); + async buildSessionDataAfterMessage(context) { + const { state, sessionId } = context; - state.update({ - conversationKey: this.getConversationKey(sessionId), - model: this.model, - }); + state.update({ + conversationKey: this.getConversationKey(sessionId), + model: this.model, + }); - return { - message: responseMessage, - sessionData: state.buildSessionData(sessionId), - }; - } catch (error) { - throw Errors.ollama.messageSendError(error); - } + return state.buildSessionData(sessionId); } /** diff --git a/models/providers/openai-responses.js b/models/providers/openai-responses.js index 63d875c..15bcd60 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), @@ -23,77 +22,65 @@ class OpenAIResponsesProvider extends BaseModel { 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 { sessionId, message, sessionData } = options; - - if (!sessionId) { - throw Errors.openAI.missingSessionId(); - } - - 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.' + ); } - await this.ensureSystemMessage(sessionId, systemInstruction); - await this.appendMessage(sessionId, 'user', message); + const conversation = await this.createConversation(); + conversationId = conversation.id; + } - const response = await this.getClient().responses.create({ - model: this.model, - conversation: conversationId, - instructions: systemInstruction, - input: [ - { - role: 'user', - content: message, - }, - ], - store: true, - }); + const response = await this.getClient().responses.create({ + model: this.model, + conversation: conversationId, + instructions: systemInstruction, + input: [ + { + role: 'user', + content: message, + }, + ], + store: true, + }); - if (response?.error) { - throw Errors.openAI.responseError(response.error.message); - } + if (response?.error) { + throw Errors.openAI.responseError(response.error.message); + } - const responseMessage = this.extractResponseText(response); + context.conversationId = conversationId; + context.lastResponseId = response.id || state.get('lastResponseId'); - if (!responseMessage) { - throw Errors.openAI.noTextContent(); - } + return response; + } - await this.storeAssistantResponse(sessionId, responseMessage); + extractResponseMessage(response) { + return this.extractResponseText(response); + } - state.update({ - conversationId, - conversationKey: this.getConversationKey(sessionId), - systemInstruction, - lastResponseId: response.id || state.get('lastResponseId'), - }); + async buildSessionDataAfterMessage(context) { + const { state, sessionId, systemInstruction, conversationId, lastResponseId } = context; - return { - message: responseMessage, - sessionData: state.buildSessionData(conversationId), - }; - } catch (error) { - throw Errors.openAI.messageSendError(error); - } + state.update({ + conversationId, + conversationKey: this.getConversationKey(sessionId), + systemInstruction, + lastResponseId, + }); + + return state.buildSessionData(conversationId); } /** diff --git a/utils/errors.js b/utils/errors.js index 970b667..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'), From c743506a1982c833b406620fd0c2de7fe1588358 Mon Sep 17 00:00:00 2001 From: Mohamed Ahmed Date: Sun, 26 Apr 2026 14:04:57 +0200 Subject: [PATCH 10/16] refactor: streamline response message extraction in Gemini, Ollama, and OpenAI providers; remove redundant methods --- .../gemini-3.1-flash-lite-preview.js | 22 +++++------- models/providers/ollama.js | 22 +++++------- models/providers/openai-responses.js | 34 ++++++++----------- 3 files changed, 33 insertions(+), 45 deletions(-) diff --git a/models/providers/gemini-3.1-flash-lite-preview.js b/models/providers/gemini-3.1-flash-lite-preview.js index d38a72e..36112af 100644 --- a/models/providers/gemini-3.1-flash-lite-preview.js +++ b/models/providers/gemini-3.1-flash-lite-preview.js @@ -39,7 +39,15 @@ class Gemini31FlashLitePreviewProvider extends BaseModel { } extractResponseMessage(response) { - return this.extractTextFromInteraction(response); + 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 buildSessionDataAfterMessage(context) { @@ -102,18 +110,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 cdaf8ac..6c51e08 100644 --- a/models/providers/ollama.js +++ b/models/providers/ollama.js @@ -29,7 +29,15 @@ class OllamaProvider extends BaseModel { } extractResponseMessage(response) { - return this.extractAssistantMessage(response); + if (!response || typeof response !== 'object') { + return ''; + } + + const content = response.message && typeof response.message.content === 'string' + ? response.message.content.trim() + : ''; + + return content; } async buildSessionDataAfterMessage(context) { @@ -116,18 +124,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 15bcd60..42ca691 100644 --- a/models/providers/openai-responses.js +++ b/models/providers/openai-responses.js @@ -67,7 +67,21 @@ class OpenAIResponsesProvider extends BaseModel { } extractResponseMessage(response) { - return this.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'); } async buildSessionDataAfterMessage(context) { @@ -132,24 +146,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(); From 6b867eb2345668ced851c676146dee85ffc22409 Mon Sep 17 00:00:00 2001 From: Mohamed Ahmed Date: Sun, 26 Apr 2026 15:41:15 +0200 Subject: [PATCH 11/16] fix: correct response handling in extractResponseMessage method; ensure proper response validation --- models/providers/gemini-3.1-flash-lite-preview.js | 6 +++--- 1 file changed, 3 insertions(+), 3 deletions(-) diff --git a/models/providers/gemini-3.1-flash-lite-preview.js b/models/providers/gemini-3.1-flash-lite-preview.js index 36112af..eb8b346 100644 --- a/models/providers/gemini-3.1-flash-lite-preview.js +++ b/models/providers/gemini-3.1-flash-lite-preview.js @@ -39,11 +39,11 @@ class Gemini31FlashLitePreviewProvider extends BaseModel { } extractResponseMessage(response) { - if (!interaction || !Array.isArray(interaction.outputs)) { + if (!response || !Array.isArray(response.outputs)) { return ''; } - return interaction.outputs + return response.outputs .filter(output => output?.type === 'text' && typeof output.text === 'string') .map(output => output.text.trim()) .filter(Boolean) @@ -75,7 +75,7 @@ class Gemini31FlashLitePreviewProvider extends BaseModel { responseFormat: this.getEvaluationResponseFormat() }); - const responseText = this.extractTextFromInteraction(interaction); + const responseText = this.extractResponseMessage(interaction); if (!responseText) { throw Errors.gemini.noEvaluationContent(); From d71383efeb0c80d3d6ec7ee6f4b18dba27e2d8fa Mon Sep 17 00:00:00 2001 From: Mohamed Ahmed Date: Sun, 26 Apr 2026 16:04:57 +0200 Subject: [PATCH 12/16] refactor: introduce native flag for environment cache checks --- .env.example | 1 - models/providers/baseModel.js | 11 ++++++----- models/providers/ollama.js | 1 + utils/environment.js | 7 ++++++- 4 files changed, 13 insertions(+), 7 deletions(-) diff --git a/.env.example b/.env.example index ecc173a..51b637d 100644 --- a/.env.example +++ b/.env.example @@ -15,7 +15,6 @@ CONVERSATION_CACHE_ENABLED=true # Optional provider-specific overrides OPENAI_CONVERSATION_CACHE_ENABLED=true GEMINI_CONVERSATION_CACHE_ENABLED=true -OLLAMA_CONVERSATION_CACHE_ENABLED=true # ... # Ollama Configuration diff --git a/models/providers/baseModel.js b/models/providers/baseModel.js index 3708775..f20068a 100644 --- a/models/providers/baseModel.js +++ b/models/providers/baseModel.js @@ -7,6 +7,7 @@ const { redisClient } = require('../../config/redis'); class BaseModel { constructor() { this.name = 'base'; + this.native = true; this.envVar = null; this._client = null; this.conversationPrefix = 'conversations:'; @@ -73,7 +74,7 @@ class BaseModel { * @returns {Promise} */ async getConversation(sessionId) { - if (!Environment.isCacheEnabled(this.envVar)) { + if (!Environment.isCacheEnabled(this.envVar, this.native)) { return []; } @@ -101,7 +102,7 @@ class BaseModel { async appendMessage(sessionId, role, content) { const message = this.normalizeConversationMessage(role, content); - if (!message || !Environment.isCacheEnabled(this.envVar)) { + if (!message || !Environment.isCacheEnabled(this.envVar, this.native)) { return; } @@ -121,7 +122,7 @@ class BaseModel { async ensureSystemMessage(sessionId, systemInstruction) { const normalizedSystemMessage = this.normalizeConversationMessage('system', systemInstruction); - if (!normalizedSystemMessage || !Environment.isCacheEnabled(this.envVar)) { + if (!normalizedSystemMessage || !Environment.isCacheEnabled(this.envVar, this.native)) { return; } @@ -175,7 +176,7 @@ class BaseModel { * @returns {Promise} */ async buildConversationForRequest(sessionId, systemInstruction, userMessage) { - if (!Environment.isCacheEnabled(this.envVar)) { + if (!Environment.isCacheEnabled(this.envVar, this.native)) { const fallbackConversation = []; const normalizedSystemMessage = this.normalizeConversationMessage('system', systemInstruction); const normalizedUserMessage = this.normalizeConversationMessage('user', userMessage); @@ -212,7 +213,7 @@ class BaseModel { * @returns {Promise} */ async clearConversation(sessionId) { - if (!Environment.isCacheEnabled(this.envVar)) { + if (!Environment.isCacheEnabled(this.envVar, this.native)) { return; } diff --git a/models/providers/ollama.js b/models/providers/ollama.js index 6c51e08..78ab593 100644 --- a/models/providers/ollama.js +++ b/models/providers/ollama.js @@ -6,6 +6,7 @@ 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(/\/+$/, ''); diff --git a/utils/environment.js b/utils/environment.js index d7a1bad..4652560 100644 --- a/utils/environment.js +++ b/utils/environment.js @@ -16,9 +16,14 @@ const getProviderEnvVar = (envVar, suffix) => (envVar ? `${envVar}_${suffix}` : * - `CONVERSATION_CACHE_ENABLED` * * @param {string} envVar - Prefijo del proveedor. + * @param {boolean} native - Indica si el proveedor es nativo. * @returns {boolean} */ -function isCacheEnabled(envVar) { +function isCacheEnabled(envVar, native = true) { +if (!native) { + return true; +} + let rawValue; if (envVar) { From e444364747281df83a4cb1ececfb02c80336e9c9 Mon Sep 17 00:00:00 2001 From: Mohamed Ahmed Date: Sun, 3 May 2026 17:08:56 +0200 Subject: [PATCH 13/16] refactor: remove max messages handling Co-authored-by: Copilot --- models/providers/baseModel.js | 5 ----- 1 file changed, 5 deletions(-) diff --git a/models/providers/baseModel.js b/models/providers/baseModel.js index f20068a..b0f3dd1 100644 --- a/models/providers/baseModel.js +++ b/models/providers/baseModel.js @@ -11,7 +11,6 @@ class BaseModel { this.envVar = null; this._client = null; this.conversationPrefix = 'conversations:'; - this.defaultConversationMaxMessages = 60; this.defaultConversationTtlSeconds = 2629800; } @@ -107,9 +106,7 @@ class BaseModel { } const key = this.getConversationKey(sessionId); - const maxMessages = Environment.getConversationMaxMessages(this.envVar); await redisClient.rPush(key, JSON.stringify(message)); - await redisClient.lTrim(key, -maxMessages, -1); await redisClient.expire(key, this.getConversationTtlSeconds()); } @@ -127,7 +124,6 @@ class BaseModel { } const key = this.getConversationKey(sessionId); - const maxMessages = Environment.getConversationMaxMessages(this.envVar); const firstRawMessage = await redisClient.lIndex(key, 0); if (!firstRawMessage) { @@ -156,7 +152,6 @@ class BaseModel { if (normalizedFirstMessage.role !== 'system') { await redisClient.lPush(key, JSON.stringify(normalizedSystemMessage)); - await redisClient.lTrim(key, -maxMessages, -1); await redisClient.expire(key, this.getConversationTtlSeconds()); return; } From d9b6667d2660879012465aaf143d7bd11e17781e Mon Sep 17 00:00:00 2001 From: Mohamed Ahmed Date: Sun, 3 May 2026 17:25:22 +0200 Subject: [PATCH 14/16] refactor: remove conversation history max messages handling 2 --- .env.example | 1 - utils/environment.js | 33 +-------------------------------- 2 files changed, 1 insertion(+), 33 deletions(-) diff --git a/.env.example b/.env.example index 51b637d..9f835c6 100644 --- a/.env.example +++ b/.env.example @@ -6,7 +6,6 @@ DEFAULT_MODEL=openai-responses RUNNER_KEY=R2D2C3PO # Conversation History (Generic) -CONVERSATION_HISTORY_MAX_MESSAGES=60 CONVERSATION_HISTORY_TTL=2629800 # Conversation Cache Global switch diff --git a/utils/environment.js b/utils/environment.js index 4652560..27e7bd9 100644 --- a/utils/environment.js +++ b/utils/environment.js @@ -50,35 +50,4 @@ if (['0', 'false', 'no', 'off'].includes(normalized)) { return true; } -/** - * Obtiene el límite de historial de conversación para el proveedor actual. - * Prioriza la variable del proveedor y luego la global. - * - * Convención esperada: - * - `${envVar}_HISTORY_MAX_MESSAGES` - * - `CONVERSATION_HISTORY_MAX_MESSAGES` - * - * @param {string} envVar - Prefijo del proveedor. - * @returns {number} - */ -function getConversationMaxMessages(envVar) { -let rawValue; - -if (envVar) { - const providerEnvVar = getProviderEnvVar(envVar, 'HISTORY_MAX_MESSAGES'); - rawValue = process.env[providerEnvVar]; -} - -if (!rawValue) { - rawValue = process.env.CONVERSATION_HISTORY_MAX_MESSAGES; -} - -const parsed = Number.parseInt(rawValue || 60, 10); -if (!Number.isInteger(parsed) || parsed <= 0) { - return 60; -} - -return parsed; -} - -module.exports = { isCacheEnabled, getConversationMaxMessages }; \ No newline at end of file +module.exports = { isCacheEnabled }; \ No newline at end of file From 34905e0a361f053faaa23c73a0e19c3c775c577e Mon Sep 17 00:00:00 2001 From: Mohamed Ahmed Date: Sun, 3 May 2026 20:48:04 +0200 Subject: [PATCH 15/16] refactor: remove buildSessionDataAfterMessage method from BaseModel and its implementations in providers Co-authored-by: Copilot --- models/providers/baseModel.js | 34 ++----------------- .../gemini-3.1-flash-lite-preview.js | 11 +----- models/providers/ollama.js | 11 +----- models/providers/openai-responses.js | 13 +------ 4 files changed, 5 insertions(+), 64 deletions(-) diff --git a/models/providers/baseModel.js b/models/providers/baseModel.js index b0f3dd1..039eb2a 100644 --- a/models/providers/baseModel.js +++ b/models/providers/baseModel.js @@ -171,22 +171,6 @@ class BaseModel { * @returns {Promise} */ async buildConversationForRequest(sessionId, systemInstruction, userMessage) { - if (!Environment.isCacheEnabled(this.envVar, this.native)) { - const fallbackConversation = []; - const normalizedSystemMessage = this.normalizeConversationMessage('system', systemInstruction); - const normalizedUserMessage = this.normalizeConversationMessage('user', userMessage); - - if (normalizedSystemMessage) { - fallbackConversation.push(normalizedSystemMessage); - } - - if (normalizedUserMessage) { - fallbackConversation.push(normalizedUserMessage); - } - - return fallbackConversation; - } - await this.ensureSystemMessage(sessionId, systemInstruction); await this.appendMessage(sessionId, 'user', userMessage); return this.getConversation(sessionId); @@ -355,11 +339,7 @@ class BaseModel { await this.storeAssistantResponse(sessionId, responseMessage); - const updatedSessionData = await this.buildSessionDataAfterMessage( - requestContext, - response, - responseMessage - ); + const updatedSessionData = await context.state.buildSessionData(context.sessionId); return { message: responseMessage, @@ -382,7 +362,7 @@ class BaseModel { } /** - * 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() { @@ -420,16 +400,6 @@ class BaseModel { throw new Error('Method extractResponseMessage must be implemented by subclasses'); } - /** - * Construye el sessionData que se devolverá al servicio de sesión. - * @param {Object} context - Contexto de la solicitud - * @param {Object} response - Respuesta cruda del proveedor - * @param {string} responseMessage - Mensaje final del asistente - * @returns {Object} - */ - async buildSessionDataAfterMessage(context, response, responseMessage) { - return context.state.buildSessionData(context.sessionId); - } } 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 eb8b346..3976101 100644 --- a/models/providers/gemini-3.1-flash-lite-preview.js +++ b/models/providers/gemini-3.1-flash-lite-preview.js @@ -50,16 +50,7 @@ class Gemini31FlashLitePreviewProvider extends BaseModel { .join('\n\n'); } - async buildSessionDataAfterMessage(context) { - const { state, sessionId, previousInteractionId } = context; - - state.update({ - previousInteractionId, - conversationKey: this.getConversationKey(sessionId), - }); - - return state.buildSessionData(previousInteractionId); - } + /** * Realiza la llamada al API de Gemini y devuelve la evaluación estructurada. diff --git a/models/providers/ollama.js b/models/providers/ollama.js index 78ab593..547602c 100644 --- a/models/providers/ollama.js +++ b/models/providers/ollama.js @@ -41,16 +41,7 @@ class OllamaProvider extends BaseModel { return content; } - async buildSessionDataAfterMessage(context) { - const { state, sessionId } = context; - - state.update({ - conversationKey: this.getConversationKey(sessionId), - model: this.model, - }); - - return state.buildSessionData(sessionId); - } + /** * Realiza la llamada al API de Ollama y devuelve la evaluación estructurada. diff --git a/models/providers/openai-responses.js b/models/providers/openai-responses.js index 42ca691..399091e 100644 --- a/models/providers/openai-responses.js +++ b/models/providers/openai-responses.js @@ -84,18 +84,7 @@ class OpenAIResponsesProvider extends BaseModel { .join('\n\n'); } - async buildSessionDataAfterMessage(context) { - const { state, sessionId, systemInstruction, conversationId, lastResponseId } = context; - - state.update({ - conversationId, - conversationKey: this.getConversationKey(sessionId), - systemInstruction, - lastResponseId, - }); - - return state.buildSessionData(conversationId); - } + /** * Realiza la llamada al API de OpenAI y devuelve la evaluación estructurada. From 224b84761cbc337e755e9ae00bd113c12793a068 Mon Sep 17 00:00:00 2001 From: Mohamed Ahmed Date: Sun, 3 May 2026 20:56:54 +0200 Subject: [PATCH 16/16] fix: update session data retrieval to use requestContext in appendMessage method Co-authored-by: Copilot --- models/providers/baseModel.js | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/models/providers/baseModel.js b/models/providers/baseModel.js index 039eb2a..4a60669 100644 --- a/models/providers/baseModel.js +++ b/models/providers/baseModel.js @@ -339,7 +339,7 @@ class BaseModel { await this.storeAssistantResponse(sessionId, responseMessage); - const updatedSessionData = await context.state.buildSessionData(context.sessionId); + const updatedSessionData = await requestContext.state.buildSessionData(requestContext.sessionId); return { message: responseMessage,