USER
User: //const { gpt, gptweb } = require("./API/gpt");
//const { dalle } = require("./API/dalle");
const { gpt4o } = require("c:/PrivateAPI/API/gpt4o");
//const { gpt3 } = require("./API/gpt3.5");
//const { prompts } = require("./API/gpt-prompts");
const { blackbox } = require("c:/PrivateAPI/API/blackbox");
const { qwen } = require("c:/PrivateAPI/API/qwen");
const { gemini } = require("c:/PrivateAPI/API/gemini");
const WebSocket = require('ws');
const Bottleneck = require('bottleneck');
const { stats, scheduleDailyReset, resetStatistics } = require('./stats'); // Импортируем модуль статистики
const fs = require('fs');
const path = require('path');
const Ajv = require('ajv');
const axios = require('axios');
const https = require('https');
const redis = require('redis');
const { promisify } = require('util');
const express = require('express');
const cors = require('cors'); // Импортируем cors
// Обработка необработанных исключений (помогут выявить обшибки, которые не обрабатываются блоками try/catch)
process.on('uncaughtException', (err) => {
console.error('Необработанное исключение:', err);
});
process.on('unhandledRejection', (reason, promise) => {
console.error('Необработанное отклонение промиса:', reason);
});
// Текущее количество подключенных клиентов
let currentConnectedClients = 0;
// Текущее количество обрабатываемых запросов
let currentGPT4oProcessingRequests = 0;
let currentGeminiProcessingRequests = 0;
let currentQwenProcessingRequests = 0;
let currentBlackboxProcessingRequests = 0;
// Объект для хранения количества подключений по IP
const ipConnectionCount = {};
// Максимальное количество подключений с одного IP
const MAX_CONNECTIONS_PER_IP = 100;
// Создание и настройка клиента Redis
const redisClient = redis.createClient();
// Подключаем клиента Redis
redisClient.connect().catch(err => {
console.error('Ошибка подключения к Redis:', err);
});
redisClient.on('error', (err) => {
console.error('Ошибка Redis:', err);
});
// Промис-обертка для команды get
const getAsync = async (key) => {
if (!redisClient.isOpen) { // Проверяем, открыт ли клиент
console.error('Клиент Redis не подключён, невозможно выполнить get:', key);
return null; // Возвращаем null, если клиент не подключен
}
try {
return await redisClient.get(key);
} catch (error) {
console.error('Ошибка при получении из Redis:', error);
throw error; //throw error чтобы выше поймали
}
};
// SSL сертификаты
const serverOptions = {
key: fs.readFileSync('C:\\OSPanel\\userdata\\config\\cert_files\\gpt.loc-key.pem'), // Приватный ключ
cert: fs.readFileSync('C:\\OSPanel\\userdata\\config\\cert_files\\gpt.loc.pem'), // Сертификат
};
// HTTPS сервер
const httpsServer = https.createServer(serverOptions);
// Создайте WebSocket-сервер, используя существующий HTTPS сервер
const wss = new WebSocket.Server({ server: httpsServer });
// Запустите сервер
httpsServer.listen(3000, () => {
console.log('HTTPS сервер и WebSocket сервер запущены на порту 3000');
});
// Создаем экземпляр Ajv (проверка входящей переписки на соответствие схеме)
const ajv = new Ajv();
// Объект для хранения действий и состояния каждого клиента
const clientActions = {};
// Объект для хранения ограничителей по каждому клиенту
const limiters = {};
// Домен
const allowedOrigin = 'https://gpt.loc';
wss.on('connection', (ws, req) => {
//console.log('Новое соединение:', req.socket.remoteAddress);
const urlParams = new URLSearchParams(req.url.split('?')[1]);
const userId = urlParams.get('token'); // Получаем UUID из URL
const userIP = req.socket.remoteAddress;
//console.log("IP-адрес - ", userIP, "Время подключения - ", new Date().toISOString());
const userHeaders = req.headers;
// Проверяем Origin
if (!checkOrigin(userHeaders.origin, ws)) {
return; // Если отклонено, выходим из обработчика
}
// Если количество подключений с этого IP превышает лимит, закрываем соединение
if (ipConnectionCount[userIP] >= MAX_CONNECTIONS_PER_IP) {
console.log(`Превышен лимит подключений для IP: ${userIP}, отказано в соединении.`);
stats.totalIPBlocks++;
ws.close(); // Закрываем соединение
return; // Выходим из обработчика
}
// Увеличиваем счетчик подключений
currentConnectedClients++;
// Обновляем максимальное количество клиентов, если текущее количество больше
if (currentConnectedClients > stats.maxConnectedClients) {
stats.maxConnectedClients = currentConnectedClients;
console.log(`Максимальное количество подключенных клиентов: ${stats.maxConnectedClients}`);
}
// Увеличиваем счетчик подключений для этого IP
ipConnectionCount[userIP] = (ipConnectionCount[userIP] || 0) + 1;
// Увеличиваем общий счетчик подключений
stats.totalConnections++;
console.log('Client connected with user ID:', userId, ' and IP:', userIP);
// Инициализируем состояние клиента, если его еще нет
if (!clientActions[userId]) {
clientActions[userId] = {
actions: [], // Массив для хранения действий
connected: true, // Статус подключения
lastRequestTime: new Date(), // Добавляем время последнего запроса
ws: ws,
processing: false;
delay: 1000, // Задержка по умолчанию, например, 1000 мс
responseTimes: [] // массив для хранения временных меток (начало и конец 3-х ответов)
};
}
else {
// обновляем ws при переподключении
clientActions[userId].connected = true;
clientActions[userId].ws = ws;
// вы можете обновить другие параметры здесь
}
if (!limiters[userId]) { console.log('Новый Bottleneck, ', limiters);
// Создаем экземпляр Bottleneck для каждого клиента (ограничение на количество запросов по Id)
limiters[userId] = new Bottleneck({
maxConcurrent: 1, // Максимальное количество одновременных задач
minTime: 1000, // Минимальное время между запросами в миллисекундах
highWater: 0 // Не разрешать в очереди больше 0 запросов (игнорировать)
});
}
ws.on('message', (message) => {
// Обновляем время последнего запроса
clientActions[userId].lastRequestTime = new Date();
// Проверяем, является ли сообщение экземпляром Buffer
if (Buffer.isBuffer(message))
message = message.toString();
try {
// Преобразуем строку JSON обратно в объект
message = JSON.parse(message);
//Получаем историю чата из запроса
let messages = message.messages;
let model = message.model;
let rewrite = message.rewrite;
// Выполняем валидацию сообщений
messages = validateMessages(messages);
// Проверяем, остались ли сообщения после валидации
if (messages.length === 0) {
const response = { code: 2, message: "Нет валидных сообщений для обработки." };
sendMessage(ws, response);
return;
}
// Урезаем последнее сообщение до 2000 символов, если оно от клиента
if (messages[messages.length - 1].role === 'user') {
messages[messages.length - 1].content = messages[messages.length - 1].content.substring(0, 2000000);
}
// Увеличиваем общий счетчик запросов
stats.totalRequests++;
if(clientActions[userId].processing){
console.lor('Игнор');
return;
}
clientActions[userId].processing = true;
// Задержка перед отправкой запроса
const responseTimes = clientActions[userId].responseTimes; // Получаем массив responseTimes
// Проверяем, прошло ли больше 15 секунд с момента окончания последнего ответа
if (responseTimes.length === 3) {
const lastEndTime = responseTimes[responseTimes.length - 1].end;
const currentTime = new Date();
// Вычисляем разницу во времени
const timeSinceLastEnd = currentTime - lastEndTime; // Время в миллисекундах
if (timeSinceLastEnd > 20000) { // Если прошло больше 15 секунд
clientActions[userId].delay = 1000; // Устанавливаем delay на 1 мс
}
}
console.log('Задержка - ', clientActions[userId].delay)
// Используем Bottleneck для управления количеством запросов
limiters[userId].schedule(() => {
// Отправка запроса на сервер чата с задержкой
setTimeout(() => {
call(messages, model, ws, userId, rewrite);
}, clientActions[userId].delay);
}).catch(error => {
stats.totalIPBlocks++;
//console.error('Rate limit exceeded, message ignored:', error);
//const response = { code: 2, message: "Превышен лимит запросов. Пожалуйста, попробуйте позже." };
// Отправка ответа клиенту
//sendMessage(ws, response);
});
} catch (error) {
stats.totalWsMessErrors++;
console.error('Error parsing JSON:', error);
const response = { code: 2, message: "Произошла ошибка. Попробуйте еще раз."};
console.log(response);
// Отправка ответа клиенту
sendMessage(ws, response);
}
});
ws.on('close', () => {
try {
// Проверка, существует ли userId в clientActions
if (clientActions[userId]) {
clientActions[userId].connected = false; // Обновление статуса подключения
delete clientActions[userId]; // Удаление состояния клиента
}
// Проверка, существует ли userId в limiters
if (limiters[userId]) {
delete limiters[userId]; // Удаление ограничителя для данного клиента
}
console.log('Client ',userId, ' disconnected');
// Уменьшаем счетчик подключений для этого IP
if (ipConnectionCount[userIP] !== undefined) {
ipConnectionCount[userIP]--;
// Если счетчик равен 0, удаляем запись из объекта
if (ipConnectionCount[userIP] === 0) {
delete ipConnectionCount[userIP];
}
}
// Уменьшаем счетчик при отключении клиента
currentConnectedClients = Math.max(currentConnectedClients - 1, 0);
} catch (error) {
console.error('Error handling client disconnect:', error);
}
});
});
// Запускаем планировщик для статистики
scheduleDailyReset();
// Таймер для проверки клиентов каждые 10 минут
setInterval(() => {
const currentTime = new Date();
for (const userId in clientActions) {
const client = clientActions[userId];
const lastRequestTime = client.lastRequestTime;
// Проверяем, прошло ли более 30 минут с последнего сообщения
if (lastRequestTime && (currentTime - lastRequestTime) > 30 * 60 * 1000) {
const response = { code: 3, message: "" };
// Отправляем сообщение клиенту
if (client.ws) {
sendMessage(client.ws, response);
// Закрываем соединение
client.ws.close();
}
}
}
console.log('Количество активных клиентов - ', Object.keys(clientActions).length);
}, 10 * 60 * 1000); // Проверяем каждые 10 минут
// Очистка старых соединений раз в минуту (по желанию)
/*setInterval(() => {
Object.keys(clientActions).forEach(userId => {
if (!clientActions[userId].connected) {
delete clientActions[userId]; // Удаляем пользователя из памяти, если он отключен
delete limiters[userId];
}
});
console.log('Количество активных клиентов - ', Object.keys(clientActions).length);
}, 10 * 60 * 1000);*/ // Очистка каждые 10 минут
// Код выводит текущие метрики использования памяти приложения
setInterval(() => {
const memoryUsage = process.memoryUsage();
console.log(`Текущие метрики использования памяти:
RSS: ${memoryUsage.rss / (1024 * 1024)} MB,
Heap Total: ${memoryUsage.heapTotal / (1024 * 1024)} MB,
Heap Used: ${memoryUsage.heapUsed / (1024 * 1024)} MB,
External: ${memoryUsage.external / (1024 * 1024)} MB`);
}, 60 * 60 * 1000); // Вывод каждый 60 секунд
// Проверяет с какого домена пришел запрос
function checkOrigin(origin, socket) {
if (origin !== allowedOrigin) {
console.log('Недопустимый домен. Соединение разорвано: ', origin);
socket.close(); // Закрываем соединение, если домен не разрешен
return false; // Возвращаем false для индикатора отклонения
}
return true; // Возвращаем true, если домен разрешен
}
// Отправляет ответ клиенту
function sendMessage(ws, message) {
if (ws.readyState === WebSocket.OPEN) {
ws.send(JSON.stringify(message));
} else {
console.warn(`Соединение закрыто. Сообщение не отправлено.`);
}
}
// Удаление из массива сообщений не соответствующих схеме
function validateMessages(messages) {
// Определяем схему для сообщений
const schema = {
type: 'object',
properties: {
role: { type: 'string' },
content: { type: 'string' }
},
required: ['role', 'content'],
additionalProperties: false
};
// Фильтруем массив, оставляя только корректные сообщения
return messages.filter(message => {
// Проверяем, соответствует ли сообщение схеме
if (ajv.validate(schema, message)) {
// Обрезаем содержимое, если оно длиннее 2000 символов
if (message.content.length > 2000000) {
message.content = message.content.substring(0, 2000000);
}
return true; // Сообщение валидно, оставляем в массиве
}
return false; // Сообщение не валидно, исключаем из массива
});
}
// Слуйчайное число
function getRandomInRange(min, max) {
return Math.floor(Math.random() * (max - min + 1)) + min;
}
// Обновляем временные метки в responseTimes
function updateResponseTimes(userId, startTime, endTime) {
try {
const responseTime = { start: startTime, end: endTime };
// Получаем массив responseTimes
const responseTimes = clientActions[userId].responseTimes;
// Добавляем новую временную метку
responseTimes.push(responseTime);
// Удаляем старые временные метки, если их больше трех
if (responseTimes.length > 3) {
responseTimes.shift(); // Удаляем старую временную метку
}
// Обновляем задержку на основе средней продолжительности ответов
updateClientDelay(userId);
console.log(clientActions[userId].responseTimes);
} catch (error) {
console.error('Ошибка при обновлении временных меток:', error);
// Можно добавить дополнительные действия в случае ошибки, например отправить уведомление
}
}
// Проверяем длительность последних трех ответов и устанавливаем новое значение delay
function updateClientDelay(userId) {
try {
const responseTimes = clientActions[userId].responseTimes;
if (responseTimes.length < 3) {
clientActions[userId].delay = 1000; // Устанавливаем delay на 1000 мс
return; // Если еще нет ответов, ничего не делаем
}
// Проверяем, все ли три ответа дольше 15 секунд
const allResponsesLongerThan15s = responseTimes.every(time => {
const duration = time.end - time.start; // Вычисляем продолжительность
return duration > 15000; // Проверяем, больше ли 15 секунд
});
// Устанавливаем delay в зависимости от проверки
if (allResponsesLongerThan15s) {
clientActions[userId].delay = 20000; // Устанавливаем delay на 20000 мс
} else {
clientActions[userId].delay = 1000; // Устанавливаем delay на 1000 мс
}
} catch (error) {
console.error('Ошибка при обновлении задержки клиента:', error);
// Можно добавить дополнительные действия в случае ошибки, например установить задержку по умолчанию
clientActions[userId].delay = 1000; // Пример установки по умолчанию
}
}
// Вызов модели
async function call(messages, model, ws, userId, rewrite){
if (model == 'gpt-4o')
return callGPT4o(messages, ws, userId, rewrite);
else if (model == 'gemini-pro')
return callGemini(messages, ws, userId, rewrite);
else if (model == 'qwen')
return callQwen(messages, ws, userId, rewrite);
else if (model == 'blackbox')
return callBlackbox(messages, ws, userId, rewrite);
else
return callGPT4o(messages, ws, userId, rewrite);
}
//============== ФУНКЦИИ API МОДЕЛЕЙ ===============================================
//==================================================================================
//==================================================================================
//==================================================================================
// GPT4o
async function callGPT4o(messages, ws, userId, rewrite) {
const lastMessageContent = messages[messages.length - 1].content; // Используем только последнее сообщение
const cacheKey = `gpt4o:${userId}:${lastMessageContent}`; // Создаем ключ кэшана основе сообщения, модели и userId
// Записываем время начала ответа
const startTime = new Date();
// Пробуем получить ответ из кэша
let cachedResponse;
try {
cachedResponse = await getAsync(cacheKey);
} catch (error) {
console.error('Ошибка при получении из Redis:', error);
// Переходим к запросу по API
}
if (cachedResponse && !rewrite) {
console.log('Отправка ответа из кэша для GPT4o');
// Имитируем задержку перед отправкой ответа из кэша (например, 1 с)
setTimeout(() => {
sendMessage(ws, JSON.parse(cachedResponse));
// Записываем время окончания ответа
const endTime = new Date();
updateResponseTimes(userId, startTime, endTime);
}, 1 * 1500); // 2 секунды задержки
return;
}
// Увеличиваем счетчик запросов находящихся в обработке
currentGPT4oProcessingRequests++;
// Обновляем максимальное количество обрабатываемых запросов, если текущее больше
if (currentGPT4oProcessingRequests > stats.maxGPT4oProcessingRequests) {
stats.maxGPT4oProcessingRequests = currentGPT4oProcessingRequests;
console.log(`Новое максимальное количество одновременно обрабатываемых запросов к GPT4o: ${stats.maxGPT4oProcessingRequests}`);
}
try {//throw new Error("Намеренная ошибка для тестирования catch");
console.log('Запрос клиента к GPT4o: ', messages[messages.length - 1]);
let fullResponse = "";
await new Promise((resolve, reject) => {
gpt4o({ messages, markdown: false }, (result) => {
if (!result.finish) {
fullResponse = result.message || "";
} else {
if(fullResponse == '')
stats.totalGPT4oLimitStops++;
let responseMessage = fullResponse || 'Попробуйте написать позже или выберите другую модель.';
console.log('Ответ GPT4o:', { code: 1, message: responseMessage });
try {
// Изменение здесь: используем метод set для установки TTL (в секундах)
redisClient.set(cacheKey, JSON.stringify({ code: 1, message: responseMessage }), { EX: 10 * 60 });
} catch (error) {
console.error('Ошибка при сохранении в Redis:', error);
}
resolve({ code: 1, message: responseMessage });
stats.totalGPT4oResponse++;
sendMessage(ws, { code: 1, message: responseMessage });
// Записываем время окончания ответа
const endTime = new Date();
updateResponseTimes(userId, startTime, endTime);
}
});
});
} catch (error) {
stats.totalGPT4oMessErrors++;
console.error('Error:', error);
sendMessage(ws, { code: 2, message: "Произошла ошибка. Попробуйте еще раз." });
}finally {
// Уменьшаем счетчик независимо от результата
currentGPT4oProcessingRequests--;
}
}
Я ввел флаг processing, он становится true когда начинается обработка запроса. Но нет обратного назначения в false, когда закончена обработка. Где это сделать в таком коде?
Assistant: корректного устанавливайте значения флага `processing` в `false` после завершения обработки запроса, вам нужно добавить соответствующее присваивание `false` в блоки завершения функций, где у вас происходит обработка.
В вашем коде следует внести изменения в асинхронные функции, которые вы выполняете (например, `callGPT4o`, `callGemini`, `callQwen`, `callBlackbox`). Сначала будет хорошо добавить освобождение флага `processing` в блоки `finally`, чтобы гарантировать, что он будет сброшен независимо от того, завершается ли обработка успешно или с ошибкой.
Вот пример того, как вы можете сделать это в функции `callGPT4o`:
```javascript
async function callGPT4o(messages, ws, userId, rewrite) {
const lastMessageContent = messages[messages.length - 1].content; // Используем только последнее сообщение
const cacheKey = `gpt4o:${userId}:${lastMessageContent}`; // Создаем ключ кэша на основе сообщения, модели и userId
// Записываем время начала ответа
const startTime = new Date();
// Пробуем получить ответ из кэша
let cachedResponse;
try {
cachedResponse = await getAsync(cacheKey);
} catch (error) {
console.error('Ошибка при получении из Redis:', error);
// Переходим к запросу по API
}
if (cachedResponse && !rewrite) {
console.log('Отправка ответа из кэша для GPT4o');
// Имитируем задержку перед отправкой ответа из кэша (например, 1 с)
setTimeout(() => {
sendMessage(ws, JSON.parse(cachedResponse));
// Записываем время окончания ответа
const endTime = new Date();
updateResponseTimes(userId, startTime, endTime);
}, 1 * 1500); // 2 секунды задержки
return;
}
// Увеличиваем счетчик запросов находящихся в обработке
currentGPT4oProcessingRequests++;
// Обновляем максимальное количество обрабатываемых запросов
if (currentGPT4oProcessingRequests > stats.maxGPT4oProcessingRequests) {
stats.maxGPT4oProcessingRequests = currentGPT4oProcessingRequests;
console.log(`Новое максимальное количество одновременно обрабатываемых запросов к GPT4o: ${stats.maxGPT4oProcessingRequests}`);
}
clientActions[userId].processing = true; // Устанавливаем флаг обработки в true
try {
console.log('Запрос клиента к GPT4o: ', messages[messages.length - 1]);
let fullResponse = "";
await new Promise((resolve, reject) => {
gpt4o({ messages, markdown: false }, (result) => {
if (!result.finish) {
fullResponse = result.message || "";
} else {
if (fullResponse == '')
stats.totalGPT4oLimitStops++;
let responseMessage = fullResponse || 'Попробуйте написать позже или выберите другую модель.';
console.log('Ответ GPT4o:', { code: 1, message: responseMessage });
try {
redisClient.set(cacheKey, JSON.stringify({ code: 1, message: responseMessage }), { EX: 10 * 60 });
} catch (error) {
console.error('Ошибка при сохранении в Redis:', error);
}
resolve({ code: 1, message: responseMessage });
stats.totalGPT4oResponse++;
sendMessage(ws, { code: 1, message: responseMessage });
// Записываем время окончания ответа
const endTime = new Date();
updateResponseTimes(userId, startTime, endTime);
}
});
});
} catch (error) {
stats.totalGPT4oMessErrors++;
console.error('Error:', error);
sendMessage(ws, { code: 2, message: "Произошла ошибка. Попробуйте еще раз." });
} finally {
// Уменьшаем счетчик независимо от результата
currentGPT4oProcessingRequests--;
clientActions[userId].processing = false; // Сбрасываем флаг обработки в false
}
}
```
Таким образом, флаг `processing` устанавливается в `false`, когда обработка завершена, что позволит избежать возможных конфликтов при попытке обработки нового запроса от того же клиента.
Вы должны внести аналогичные изменения в другие функции, такие как `callGemini`, `callQwen` и `callBlackbox`, чтобы гарантировать, что флаг будет сбрасываться и там.
User: теперь processing станет false независимо от того как завершит работу функция callGPT4o?
Assistant: , совершенно верно! Используя блок `finally`, вы гарантируете, что флаг `processing` станет `false` независимо от того, как завершится функция `callGPT4o`.
Вот как это работает:
- Если обработка запроса выполнена успешно (без исключений), флаг `processing` будет установлен в `false` в блоке `finally` после завершения блока `try`, и независимо от результата обработки (успех или ошибка).
- Если во время выполнения промиса произойдет ошибка, управление перейдет в блок `catch`, где вы можете обработать ошибку, а затем в блок `finally`, где снова будет установлен `processing` в `false`.
Таким образом, блок `finally` будет выполнен в любом случае — как при успешном завершении функции, так и при возникновении ошибки, и позволит вам корректно сбросить состояние.
Таким образом, ваше использование блока `finally` позволяет избежать ситуации, когда флаг `processing` остается `true`, если происходит ошибка. Это обеспечивает более надежное управление состоянием для каждого клиента.
Если вы внесете аналогичные изменения в другие функции обработки, результаты будут аналогичными для всех аспектов вашего кода.
User: А если сработает этот блок
if (cachedResponse && !rewrite) {
console.log('Отправка ответа из кэша для GPT4o');
// Имитируем задержку перед отправкой ответа из кэша (например, 1 с)
setTimeout(() => {
sendMessage(ws, JSON.parse(cachedResponse));
// Записываем время окончания ответа
const endTime = new Date();
updateResponseTimes(userId, startTime, endTime);
}, 1 * 1500); // 2 секунды задержки
return;
}
Assistant: сработает блок:
```javascript
if (cachedResponse && !rewrite) {
console.log('Отправка ответа из кэша для GPT4o');
// Имитируем задержку перед отправкой ответа из кэша (например, 1 с)
setTimeout(() => {
sendMessage(ws, JSON.parse(cachedResponse));
// Записываем время окончания ответа
const endTime = new Date();
updateResponseTimes(userId, startTime, endTime);
}, 1 * 1500); // 2 секунды задержки
return;
}
```
то функция завершится вызовом `return`, и блок `finally` не будет выполнен, так как `finally` привязан к текущему контексту выполнения, который был бы и завершён до этого в случае, если выполнится `return`. Это значит, что флаг `processing` не будет сброшен в этом случае, и это может привести к ним к тому, что последующий запрос от того же клиента не будет обработан должным образом.
Чтобы избежать этого, вам нужно убедиться, что флаг `processing` устанавливается в `true` перед тем, как отправлять сообщение из кэша, и сбрасывается в `false` как в случае удачной обработки, так и в случае получения ответа из кэша.
Вы можете изменить структуру кода, чтобы включить установку и сброс флага `processing` также в блоке управления кэшем. Вот как можно сделать это:
```javascript
async function callGPT4o(messages, ws, userId, rewrite) {
const lastMessageContent = messages[messages.length - 1].content;
const cacheKey = `gpt4o:${userId}:${lastMessageContent}`;
const startTime = new Date();
let cachedResponse;
try {
cachedResponse = await getAsync(cacheKey);
} catch (error) {
console.error('Ошибка при получении из Redis:', error);
}
// Устанавливаем флаг обработки в true
clientActions[userId].processing = true;
if (cachedResponse && !rewrite) {
console.log('Отправка ответа из кэша для GPT4o');
setTimeout(() => {
sendMessage(ws, JSON.parse(cachedResponse));
const endTime = new Date();
updateResponseTimes(userId, startTime, endTime);
clientActions[userId].processing = false; // Сбрасываем флаг обработки здесь
}, 1500);
return; // Завершаем обработку функции
}
currentGPT4oProcessingRequests++;
if (currentGPT4oProcessingRequests > stats.maxGPT4oProcessingRequests) {
stats.maxGPT4oProcessingRequests = currentGPT4oProcessingRequests;
console.log(`Новое максимальное количество одновременно обрабатываемых запросов к GPT4o: ${stats.maxGPT4oProcessingRequests}`);
}
try {
console.log('Запрос клиента к GPT4o: ', messages[messages.length - 1]);
let fullResponse = "";
await new Promise((resolve, reject) => {
gpt4o({ messages, markdown: false }, (result) => {
if (!result.finish) {
fullResponse = result.message || "";
} else {
if (fullResponse == '')
stats.totalGPT4oLimitStops++;
let responseMessage = fullResponse || 'Попробуйте написать позже или выберите другую модель.';
console.log('Ответ GPT4o:', { code: 1, message: responseMessage });
try {
redisClient.set(cacheKey, JSON.stringify({ code: 1, message: responseMessage }), { EX: 10 * 60 });
} catch (error) {
console.error('Ошибка при сохранении в Redis:', error);
}
resolve({ code: 1, message: responseMessage });
stats.totalGPT4oResponse++;
sendMessage(ws, { code: 1, message: responseMessage });
const endTime = new Date();
updateResponseTimes(userId, startTime, endTime);
}
});
});
} catch (error) {
stats.totalGPT4oMessErrors++;
console.error('Error:', error);
sendMessage(ws, { code: 2, message: "Произошла ошибка. Попробуйте еще раз." });
} finally {
currentGPT4oProcessingRequests--;
// Флаг обработки сбрасывается в любом случае при окончании запроса
if (clientActions[userId]) {
clientActions[userId].processing = false;
}
}
}
```
С помощью этого изменения теперь флаг `processing` будет устанавливаться в `true` сразу после попытки получить кэш и сбрасываться в `false` после обработки запроса, независимо от того, был ли ответ получен из кэша или он был обработан путем вызова API. Это гарантирует корректное управление состоянием клиента и позволяет избежать проблем с повторной отправкой запросов от клиента.
User: setInterval(() => {
const currentTime = new Date();
for (const userId in clientActions) {
const client = clientActions[userId];
client.processing = false;
const lastRequestTime = client.lastRequestTime;
// Проверяем, прошло ли более 30 минут с последнего сообщения
if (lastRequestTime && (currentTime - lastRequestTime) > 30 * 60 * 1000) {
const response = { code: 3, message: "" };
// Отправляем сообщение клиенту
if (client.ws) {
sendMessage(client.ws, response);
// Закрываем соединение
client.ws.close();
}
}
}
console.log('Количество активных клиентов - ', Object.keys(clientActions).length);
}, 5 * 60 * 1000); // Проверяем каждые 10 минут
привально ли сбрасывается processing для всех клиентов
Assistant: вашем коде:
```javascript
setInterval(() => {
const currentTime = new Date();
for (const userId in clientActions) {
const client = clientActions[userId];
client.processing = false; // Сбрасываем флаг processing
const lastRequestTime = client.lastRequestTime;
// Проверяем, прошло ли более 30 минут с последнего сообщения
if (lastRequestTime && (currentTime - lastRequestTime) > 30 * 60 * 1000) {
const response = { code: 3, message: "" };
// Отправляем сообщение клиенту
if (client.ws) {
sendMessage(client.ws, response);
// Закрываем соединение
client.ws.close();
}
}
}
console.log('Количество активных клиентов - ', Object.keys(clientActions).length);
}, 5 * 60 * 1000); // Проверяем каждые 5 минут
```
Сброс флага `processing` происходит для всех клиентов каждый раз, когда срабатывает `setInterval`, то есть каждые 5 минут. Это означает, что даже если клиент все еще обрабатывает запрос (находится в состоянии `processing = true`), этот флаг будет сброшен в `false` при следующем запуске интервала.
### Потенциальные проблемы:
1. **Ложное срабатывание**: Если клиент все еще обрабатывает запрос, и в это время срабатывает интервал, то сброс флага `processing` на `false` отменит фактическое состояние обработки. Это может привести к путанице или конфликтам, если клиент попытается отправить новый запрос, в то время как фактическая обработка еще не окончена.
2. **Проблемы с идентификацией обработанных и не обработанных запросов**: Если вы сбрасываете флаг `processing` в зависимости от времени, а не от факта завершения обработки текущего запроса, это может привести к ошибкам обработки следующих запросов.
### Рекомендации:
Чтобы избежать этих проблем, вам следует обрабатывать состояние `processing` изнутри вашей логики обработки запросов, а не в интервале. Вы можете сделать сброс флага `processing` только после завершения обработки в обрабатывающих функциях, таких как `callGPT4o`, как мы обсуждали ранее, и убрать строку `client.processing = false;` из `setInterval`.
Вот как может выглядеть улучшенное решение:
- Убедитесь, что `processing` сбрасывается только в во время выполнения функции обработки.
- В `setInterval` вы можете проверить состояние клиента и отправлять уведомления или закрывать соединение только в зависимости от времени последнего запроса, без изменения флага `processing`.
Пример измененного функционала, чтобы сброс обработки происходил только после завершения работы функции:
```javascript
setInterval(() => {
const currentTime = new Date();
for (const userId in clientActions) {
const client = clientActions[userId];
const lastRequestTime = client.lastRequestTime;
// Проверяем, прошло ли более 30 минут с последнего сообщения
if (lastRequestTime && (currentTime - lastRequestTime) > 30 * 60 * 1000) {
const response = { code: 3, message: "" };
// Отправляем сообщение клиенту
if (client.ws) {
sendMessage(client.ws, response);
// Закрываем соединение
client.ws.close();
}
}
}
console.log('Количество активных клиентов - ', Object.keys(clientActions).length);
}, 5 * 60 * 1000); // Проверяем каждые 5 минут
```
Таким образом, вы обеспечите, что управление состоянием обработки остается корректным и не будет сбрасываться в ненадлежащий момент.
User: просто все ли случаи учтены когда должен быть сброшен флаг. Надо исключить случай когда флан останется true после обработки запроса клиента, иначе он не сможет больше писать
Assistant: правы, важно учесть все случаи, когда флаг `processing` должен быть сброшен. Чтобы гарантировать, что этот флаг не останется установленным в `true`, даже если клиент сделает дополнительные запросы, нужно убедиться, что его состояние корректно управляется в следующих ситуациях:
1. **После успешной обработки запроса**: Флаг `processing` должен быть сброшен, как только обработка завершена, независимо от того, был ли этот ответ получен из кэша или от API.
2. **После обработки ошибки**: Если произошла ошибка в процессе обработки запроса, флаг также должен быть сброшен, так как ошибка не должна блокировать клиента от отправки новых запросов.
3. **Закрытие соединения**: Если клиент отключается или connection closes (например, timeout срабатывает на стороне сервера), необходимо сбрасывать флаг `processing`.
4. **В период опроса**: Избегайте сброса флага `processing` в блоке `setInterval`, так как это может привести к неправильному управлению состоянием.
Вот пример, как можно полностью управлять состоянием флага `processing` в функции `callGPT4o`, учитывая все упомянутые случаи:
```javascript
async function callGPT4o(messages, ws, userId, rewrite) {
const lastMessageContent = messages[messages.length - 1].content;
const cacheKey = `gpt4o:${userId}:${lastMessageContent}`;
const startTime = new Date();
let cachedResponse;
try {
// Пробуем получить ответ из кэша
cachedResponse = await getAsync(cacheKey);
} catch (error) {
console.error('Ошибка при получении из Redis:', error);
}
// Устанавливаем флаг обработки в true
clientActions[userId].processing = true;
if (cachedResponse && !rewrite) {
console.log('Отправка ответа из кэша для GPT4o');
setTimeout(() => {
sendMessage(ws, JSON.parse(cachedResponse));
const endTime = new Date();
updateResponseTimes(userId, startTime, endTime);
clientActions[userId].processing = false; // Сбрасываем флаг обработки здесь
}, 1500);
return; // Завершаем обработку функции
}
currentGPT4oProcessingRequests++;
// Логика обработки запросов
try {
console.log('Запрос клиента к GPT4o: ', messages[messages.length - 1]);
let fullResponse = "";
await new Promise((resolve, reject) => {
gpt4o({ messages, markdown: false }, (result) => {
if (!result.finish) {
fullResponse = result.message || "";
} else {
// Обработка ответа
let responseMessage = fullResponse || 'Попробуйте написать позже или выберите другую модель.';
console.log('Ответ GPT4o:', { code: 1, message: responseMessage });
try {
redisClient.set(cacheKey, JSON.stringify({ code: 1, message: responseMessage }), { EX: 10 * 60 });
} catch (error) {
console.error('Ошибка при сохранении в Redis:', error);
}
resolve({ code: 1, message: responseMessage });
stats.totalGPT4oResponse++;
sendMessage(ws, { code: 1, message: responseMessage });
const endTime = new Date();
updateResponseTimes(userId, startTime, endTime);
}
});
});
} catch (error) {
stats.totalGPT4oMessErrors++;
console.error('Error:', error);
sendMessage(ws, { code: 2, message: "Произошла ошибка. Попробуйте еще раз." });
} finally {
currentGPT4oProcessingRequests--;
// Сброс флага обработки
if (clientActions[userId]) {
clientActions[userId].processing = false;
}
}
}
```
### Обработка отключения клиента
Дополнительным аспектом, который нужно учесть, является корректное управление состоянием `processing` при закрытии соединения. Для этого нужно убедиться, что в обработчике события `close` сбрасывается флаг:
```javascript
ws.on('close', () => {
try {
// Проверка, существует ли userId в clientActions
if (clientActions[userId]) {
clientActions[userId].connected = false; // Обновление статуса подключения
clientActions[userId].processing = false; // Сбрасываем флаг processing
delete clientActions[userId]; // Удаление состояния клиента
}
// Проверка, существует ли userId в limiters
if (limiters[userId]) {
delete limiters[userId]; // Удаление ограничителя для данного клиента
}
console.log('Client ',userId, ' disconnected');
// Уменьшаем счетчик подключений для этого IP
// ... остальной код
} catch (error) {
console.error('Error handling client disconnect:', error);
}
});
```
### Заключение
Таким образом, правильное управление состоянием флага `processing` будет гарантировать, что после завершения обработки или возникновения ошибки, клиент сможет отправить новый запрос и не столкнется с блокировкой. Применение вышеуказанных изменений обеспечит корректное и безопасное взаимодействие между клиентом и сервером.
User: ws.on('close', () => {
try {
// Проверка, существует ли userId в clientActions
if (clientActions[userId]) {
clientActions[userId].connected = false; // Обновление статуса подключения
clientActions[userId].processing = false; // Сбрасываем флаг processing
delete clientActions[userId]; // Удаление состояния клиента
}
// Проверка, существует ли userId в limiters
if (limiters[userId]) {
delete limiters[userId]; // Удаление ограничителя для данного клиента
}
console.log('Client ',userId, ' disconnected');
// Уменьшаем счетчик подключений для этого IP
// ... остальной код
} catch (error) {
console.error('Error handling client disconnect:', error);
}
});
В таком случае после перезагрузки приложение не сможет контролировать идет ли обработка предыдущего запроса. Может перенести удаление clientActions[userId] и limiters[userId] сюда
setInterval(() => {
const currentTime = new Date();
for (const userId in clientActions) {
const client = clientActions[userId];
const lastRequestTime = client.lastRequestTime;
// Проверяем, прошло ли более 30 минут с последнего сообщения
if (lastRequestTime && (currentTime - lastRequestTime) > 30 * 60 * 1000) {
const response = { code: 3, message: "" };
// Отправляем сообщение клиенту
if (client.ws) {
sendMessage(client.ws, response);
// Закрываем соединение
client.ws.close();
}
}
}
console.log('Количество активных клиентов - ', Object.keys(clientActions).length);
}, 5 * 60 * 1000); // Проверяем каждые 10 минут
Assistant: