fix potential multiple requests
This commit is contained in:
@@ -943,6 +943,7 @@ class YouProvider {
|
|||||||
const cleanup = async () => {
|
const cleanup = async () => {
|
||||||
clearTimeout(responseTimeout);
|
clearTimeout(responseTimeout);
|
||||||
clearTimeout(customEndMarkerTimer);
|
clearTimeout(customEndMarkerTimer);
|
||||||
|
clearTimeout(errorTimer);
|
||||||
if (heartbeatInterval) {
|
if (heartbeatInterval) {
|
||||||
clearInterval(heartbeatInterval);
|
clearInterval(heartbeatInterval);
|
||||||
heartbeatInterval = null;
|
heartbeatInterval = null;
|
||||||
@@ -960,22 +961,18 @@ class YouProvider {
|
|||||||
|
|
||||||
// 缓存
|
// 缓存
|
||||||
let buffer = '';
|
let buffer = '';
|
||||||
let clientPointer = 0; // 客户端指针
|
|
||||||
let proxyPointer = 0; // 反代指针
|
|
||||||
let heartbeatInterval = null; // 心跳计时器
|
let heartbeatInterval = null; // 心跳计时器
|
||||||
|
let errorTimer = null; // 错误计时器
|
||||||
|
const ERROR_TIMEOUT = 15000;
|
||||||
const self = this;
|
const self = this;
|
||||||
page.exposeFunction("callback" + traceId, async (event, data) => {
|
page.exposeFunction("callback" + traceId, async (event, data) => {
|
||||||
if (isEnding) return;
|
if (isEnding) return;
|
||||||
|
|
||||||
switch (event) {
|
switch (event) {
|
||||||
case "resetProxyPointer":
|
|
||||||
proxyPointer = 0; // 重置反代指针
|
|
||||||
break;
|
|
||||||
case "youChatToken": {
|
case "youChatToken": {
|
||||||
data = JSON.parse(data);
|
data = JSON.parse(data);
|
||||||
let tokenContent = data.youChatToken;
|
let tokenContent = data.youChatToken;
|
||||||
buffer += tokenContent;
|
buffer += tokenContent;
|
||||||
proxyPointer += tokenContent.length; // 每次接收到数据,反代指针增加
|
|
||||||
|
|
||||||
if (buffer.endsWith('\\') && !buffer.endsWith('\\\\')) {
|
if (buffer.endsWith('\\') && !buffer.endsWith('\\\\')) {
|
||||||
// 等待下一个字符
|
// 等待下一个字符
|
||||||
@@ -1001,6 +998,12 @@ class YouProvider {
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// 重置错误计时器
|
||||||
|
if (errorTimer) {
|
||||||
|
clearTimeout(errorTimer);
|
||||||
|
errorTimer = null;
|
||||||
|
}
|
||||||
|
|
||||||
// 检测 'unusual query volume'
|
// 检测 'unusual query volume'
|
||||||
if (processedContent.includes('unusual query volume')) {
|
if (processedContent.includes('unusual query volume')) {
|
||||||
const warningMessage = "您在 you.com 账号的使用已达上限,当前(default/agent)模式已进入冷却期(CD)。请切换模式(default/agent[custom])或耐心等待冷却结束后再继续使用。";
|
const warningMessage = "您在 you.com 账号的使用已达上限,当前(default/agent)模式已进入冷却期(CD)。请切换模式(default/agent[custom])或耐心等待冷却结束后再继续使用。";
|
||||||
@@ -1033,21 +1036,17 @@ class YouProvider {
|
|||||||
break;
|
break;
|
||||||
}
|
}
|
||||||
|
|
||||||
if (proxyPointer >= clientPointer) {
|
process.stdout.write(processedContent);
|
||||||
let newContent = processedContent.slice(clientPointer - proxyPointer);
|
accumulatedResponse += processedContent;
|
||||||
clientPointer = proxyPointer; // 更新客户端指针
|
|
||||||
|
|
||||||
process.stdout.write(newContent);
|
|
||||||
accumulatedResponse += newContent;
|
|
||||||
|
|
||||||
if (Date.now() - startTime >= 20000) {
|
if (Date.now() - startTime >= 20000) {
|
||||||
responseAfter20Seconds += newContent;
|
responseAfter20Seconds += processedContent;
|
||||||
}
|
}
|
||||||
|
|
||||||
if (stream) {
|
if (stream) {
|
||||||
emitter.emit("completion", traceId, newContent);
|
emitter.emit("completion", traceId, processedContent);
|
||||||
} else {
|
} else {
|
||||||
finalResponse += newContent;
|
finalResponse += processedContent;
|
||||||
}
|
}
|
||||||
|
|
||||||
// 检查自定义结束标记
|
// 检查自定义结束标记
|
||||||
@@ -1067,7 +1066,6 @@ class YouProvider {
|
|||||||
unusualQueryVolume: unusualQueryVolumeTriggered,
|
unusualQueryVolume: unusualQueryVolumeTriggered,
|
||||||
});
|
});
|
||||||
}
|
}
|
||||||
}
|
|
||||||
break;
|
break;
|
||||||
}
|
}
|
||||||
case "customEndMarkerEnabled":
|
case "customEndMarkerEnabled":
|
||||||
@@ -1091,15 +1089,17 @@ class YouProvider {
|
|||||||
case "error": {
|
case "error": {
|
||||||
if (isEnding) return; // 如果已经结束,则忽略错误
|
if (isEnding) return; // 如果已经结束,则忽略错误
|
||||||
console.error("请求发生错误", data);
|
console.error("请求发生错误", data);
|
||||||
const errorMessage = data.message || "未知错误";
|
if (errorTimer) {
|
||||||
const clientErrorMessage = `请求发生错误: ${errorMessage} (错误详情已记录到日志中)`;
|
clearTimeout(errorTimer);
|
||||||
emitter.emit("completion", traceId, clientErrorMessage);
|
}
|
||||||
|
errorTimer = setTimeout(async () => {
|
||||||
|
console.log("连接超时,终止请求");
|
||||||
|
const errorMessage = "连接中断,未收到服务器响应";
|
||||||
|
emitter.emit("completion", traceId, errorMessage);
|
||||||
finalResponse += ` (${errorMessage})`;
|
finalResponse += ` (${errorMessage})`;
|
||||||
isEnding = true;
|
isEnding = true;
|
||||||
setTimeout(async () => {
|
|
||||||
await cleanup();
|
await cleanup();
|
||||||
emitter.emit("end", traceId);
|
emitter.emit("end", traceId);
|
||||||
}, 1000);
|
|
||||||
self.logger.logRequest({
|
self.logger.logRequest({
|
||||||
email: username,
|
email: username,
|
||||||
time: requestTime,
|
time: requestTime,
|
||||||
@@ -1108,11 +1108,13 @@ class YouProvider {
|
|||||||
completed: false,
|
completed: false,
|
||||||
unusualQueryVolume: unusualQueryVolumeTriggered,
|
unusualQueryVolume: unusualQueryVolumeTriggered,
|
||||||
});
|
});
|
||||||
|
}, ERROR_TIMEOUT);
|
||||||
break;
|
break;
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
});
|
});
|
||||||
|
|
||||||
|
|
||||||
// proxy response
|
// proxy response
|
||||||
const req_param = new URLSearchParams();
|
const req_param = new URLSearchParams();
|
||||||
req_param.append("page", "1");
|
req_param.append("page", "1");
|
||||||
@@ -1261,34 +1263,18 @@ class YouProvider {
|
|||||||
let evtSource;
|
let evtSource;
|
||||||
const callbackName = "callback" + traceId;
|
const callbackName = "callback" + traceId;
|
||||||
let isEnding = false;
|
let isEnding = false;
|
||||||
let retryCount = 0;
|
|
||||||
const maxRetries = 5;
|
|
||||||
let customEndMarkerTimer = null;
|
let customEndMarkerTimer = null;
|
||||||
|
|
||||||
function connect() {
|
function connect() {
|
||||||
// 重置反代指针
|
|
||||||
window[callbackName]("resetProxyPointer", "");
|
|
||||||
|
|
||||||
evtSource = new EventSource(url);
|
evtSource = new EventSource(url);
|
||||||
|
|
||||||
evtSource.onerror = (error) => {
|
evtSource.onerror = (error) => {
|
||||||
if (isEnding) return;
|
if (isEnding) return;
|
||||||
|
|
||||||
// 检查是否需要重连
|
|
||||||
if (retryCount < maxRetries) {
|
|
||||||
retryCount++;
|
|
||||||
setTimeout(() => {
|
|
||||||
console.log(`连接断开,尝试重新连接 (${retryCount}/${maxRetries})`);
|
|
||||||
connect();
|
|
||||||
}, 500 * retryCount); // 指数退避
|
|
||||||
} else {
|
|
||||||
window[callbackName]("error", error);
|
window[callbackName]("error", error);
|
||||||
}
|
|
||||||
};
|
};
|
||||||
|
|
||||||
evtSource.addEventListener("youChatToken", (event) => {
|
evtSource.addEventListener("youChatToken", (event) => {
|
||||||
if (isEnding) return;
|
if (isEnding) return;
|
||||||
retryCount = 0; // 重置重试计数器
|
|
||||||
const data = JSON.parse(event.data);
|
const data = JSON.parse(event.data);
|
||||||
window[callbackName]("youChatToken", JSON.stringify(data));
|
window[callbackName]("youChatToken", JSON.stringify(data));
|
||||||
|
|
||||||
@@ -1308,7 +1294,6 @@ class YouProvider {
|
|||||||
|
|
||||||
evtSource.onmessage = (event) => {
|
evtSource.onmessage = (event) => {
|
||||||
if (isEnding) return;
|
if (isEnding) return;
|
||||||
retryCount = 0;
|
|
||||||
const data = JSON.parse(event.data);
|
const data = JSON.parse(event.data);
|
||||||
if (data.youChatToken) {
|
if (data.youChatToken) {
|
||||||
window[callbackName]("youChatToken", JSON.stringify(data));
|
window[callbackName]("youChatToken", JSON.stringify(data));
|
||||||
|
|||||||
Reference in New Issue
Block a user