File size: 14,871 Bytes
ddce7e8
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
212
213
214
215
216
217
218
219
220
221
222
223
224
225
226
227
228
229
230
231
232
233
234
235
236
237
238
239
240
241
242
243
244
245
246
247
248
249
250
251
252
253
254
255
256
257
258
259
260
261
262
263
264
265
266
267
268
269
270
271
272
273
274
275
276
277
278
279
280
281
282
283
284
285
286
287
288
289
290
291
292
293
294
295
296
297
298
299
300
301
302
303
304
305
306
307
308
309
310
311
312
313
314
315
316
317
318
319
320
321
322
323
324
325
326
327
328
329
330
331
332
333
334
335
336
337
338
339
340
341
342
343
344
345
346
347
348
349
350
351
352
353
354
355
356
357
358
359
360
361
362
363
364
365
366
367
368
369
370
371
372
373
374
375
376
377
378
379
380
381
382
383
384
385
386
387
388
389
390
391
392
393
394
395
396
397
398
399
400
401
402
403
404
405
406
407
408
409
410
411
412
413
414
415
416
/**
 * SSE 流式响应和心跳机制工具模块
 * 提供统一的流式响应处理、心跳保活、429/503重试等功能
 */

import config from '../config/config.js';
import logger from '../utils/logger.js';
import memoryManager, { registerMemoryPoolCleanup } from '../utils/memoryManager.js';
import { DEFAULT_HEARTBEAT_INTERVAL, LONG_COOLDOWN_THRESHOLD } from '../constants/index.js';
import tokenCooldownManager from '../auth/token_cooldown_manager.js';
import quotaManager from '../auth/quota_manager.js';
import { getGroupKey } from '../utils/modelGroups.js';

// ==================== 心跳机制(防止 CF 超时) ====================
const HEARTBEAT_INTERVAL = config.server.heartbeatInterval || DEFAULT_HEARTBEAT_INTERVAL;
const SSE_HEARTBEAT = Buffer.from(': heartbeat\n\n');

/**
 * 创建心跳定时器
 * @param {Response} res - Express响应对象
 * @returns {NodeJS.Timeout} 定时器
 */
export const createHeartbeat = (res) => {
  const timer = setInterval(() => {
    if (!res.writableEnded) {
      res.write(SSE_HEARTBEAT);
    } else {
      clearInterval(timer);
    }
  }, HEARTBEAT_INTERVAL);

  // 响应结束时清理
  res.on('close', () => clearInterval(timer));
  res.on('finish', () => clearInterval(timer));

  return timer;
};

// ==================== 预编译的常量字符串(避免重复创建) ====================
const SSE_PREFIX = Buffer.from('data: ');
const SSE_SUFFIX = Buffer.from('\n\n');
const SSE_DONE = Buffer.from('data: [DONE]\n\n');

/**
 * 生成响应元数据
 * @returns {{id: string, created: number}}
 */
export const createResponseMeta = () => ({
  id: `chatcmpl-${Date.now()}`,
  created: Math.floor(Date.now() / 1000)
});

/**
 * 设置流式响应头
 * @param {Response} res - Express响应对象
 */
export const setStreamHeaders = (res) => {
  res.setHeader('Content-Type', 'text/event-stream');
  res.setHeader('Cache-Control', 'no-cache');
  res.setHeader('Connection', 'keep-alive');
  res.setHeader('X-Accel-Buffering', 'no'); // 禁用 nginx 缓冲
  // 立即发送响应头,确保客户端尽快建立连接
  res.flushHeaders();
};

// ==================== 对象池(减少 GC) ====================
const chunkPool = [];

/**
 * 从对象池获取 chunk 对象
 * @returns {Object}
 */
export const getChunkObject = () => chunkPool.pop() || { choices: [{ index: 0, delta: {}, finish_reason: null }] };

/**
 * 释放 chunk 对象回对象池
 * @param {Object} obj 
 */
export const releaseChunkObject = (obj) => {
  const maxSize = memoryManager.getPoolSizes().chunk;
  if (chunkPool.length < maxSize) chunkPool.push(obj);
};

// 注册内存清理回调
registerMemoryPoolCleanup(chunkPool, () => memoryManager.getPoolSizes().chunk);

/**
 * 获取当前对象池大小(用于监控)
 * @returns {number}
 */
export const getChunkPoolSize = () => chunkPool.length;

/**
 * 清空对象池
 */
export const clearChunkPool = () => {
  chunkPool.length = 0;
};

/**
 * 零拷贝写入流式数据
 * @param {Response} res - Express响应对象
 * @param {Object} data - 要发送的数据
 */
export const writeStreamData = (res, data) => {
  const json = JSON.stringify(data);
  res.write(SSE_PREFIX);
  res.write(json);
  res.write(SSE_SUFFIX);
  // 立即刷新缓冲区,确保数据实时发送给客户端
  if (typeof res.flush === 'function') {
    res.flush();
  }
};

/**
 * 结束流式响应
 * @param {Response} res - Express响应对象
 */
export const endStream = (res, isWriteDone = true) => {
  if (res.writableEnded) return;
  if (isWriteDone) res.write(SSE_DONE);
  res.end();
};

// ==================== 通用重试工具(处理 429/503) ====================

function sleep(ms) {
  return new Promise(resolve => setTimeout(resolve, ms));
}

function parseDurationToMs(value) {
  if (value === null || value === undefined) return null;
  if (typeof value === 'number' && Number.isFinite(value)) return Math.max(0, Math.floor(value));
  if (typeof value !== 'string') return null;

  const s = value.trim();
  if (!s) return null;

  // e.g. "295.285334ms"
  const msMatch = s.match(/^(\d+(\.\d+)?)\s*ms$/i);
  if (msMatch) return Math.max(0, Math.floor(Number(msMatch[1])));

  // e.g. "0.295285334s"
  const secMatch = s.match(/^(\d+(\.\d+)?)\s*s$/i);
  if (secMatch) return Math.max(0, Math.floor(Number(secMatch[1]) * 1000));

  // plain number in string: treat as ms
  const num = Number(s);
  if (Number.isFinite(num)) return Math.max(0, Math.floor(num));
  return null;
}

function tryParseJson(value) {
  if (!value) return null;
  if (typeof value === 'object') return value;
  if (typeof value !== 'string') return null;
  try {
    return JSON.parse(value);
  } catch {
    // Some messages embed JSON inside a string; try to salvage a JSON object substring.
    const first = value.indexOf('{');
    const last = value.lastIndexOf('}');
    if (first !== -1 && last !== -1 && last > first) {
      const sliced = value.slice(first, last + 1);
      try {
        return JSON.parse(sliced);
      } catch { }
    }
    return null;
  }
}

function extractUpstreamErrorBody(error) {
  // UpstreamApiError created by createApiError(...) stores rawBody
  if (error?.isUpstreamApiError && error.rawBody) {
    return tryParseJson(error.rawBody) || error.rawBody;
  }
  // axios-like error
  if (error?.response?.data) {
    return tryParseJson(error.response.data) || error.response.data;
  }
  // fallback: try parse message
  return tryParseJson(error?.message);
}

function getUpstreamRetryDelayMs(error) {
  // Prefer explicit hints from upstream payload (RetryInfo/quotaResetDelay/quotaResetTimeStamp)
  const body = extractUpstreamErrorBody(error);
  const root = (body && typeof body === 'object') ? body : null;
  const inner = root?.error || root;
  const details = Array.isArray(inner?.details) ? inner.details : [];

  let bestMs = null;
  for (const d of details) {
    if (!d || typeof d !== 'object') continue;

    // google.rpc.RetryInfo: { retryDelay: "0.295285334s" }
    const retryDelayMs = parseDurationToMs(d.retryDelay);
    if (retryDelayMs !== null) bestMs = bestMs === null ? retryDelayMs : Math.max(bestMs, retryDelayMs);

    // google.rpc.ErrorInfo metadata: { quotaResetDelay: "295.285334ms", quotaResetTimeStamp: "..." }
    const meta = d.metadata && typeof d.metadata === 'object' ? d.metadata : null;
    const quotaResetDelayMs = parseDurationToMs(meta?.quotaResetDelay);
    if (quotaResetDelayMs !== null) bestMs = bestMs === null ? quotaResetDelayMs : Math.max(bestMs, quotaResetDelayMs);

    const ts = meta?.quotaResetTimeStamp;
    if (typeof ts === 'string') {
      const t = Date.parse(ts);
      if (Number.isFinite(t)) {
        const deltaMs = Math.max(0, t - Date.now());
        bestMs = bestMs === null ? deltaMs : Math.max(bestMs, deltaMs);
      }
    }
  }

  // If it's the capacity exhausted case, still retry but avoid hammering.
  const reason = details.find(d => d?.reason)?.reason;
  if (reason === 'MODEL_CAPACITY_EXHAUSTED') {
    bestMs = bestMs === null ? 1000 : Math.max(bestMs, 1000);
  }

  return bestMs;
}

function computeBackoffMs(attempt, explicitDelayMs) {
  // attempt starts from 0 for first call; on first retry attempt=1
  const maxMs = 20_000;
  const hasExplicit = Number.isFinite(explicitDelayMs) && explicitDelayMs !== null;
  const baseMs = hasExplicit ? Math.max(0, Math.floor(explicitDelayMs)) : 500;
  const exp = Math.min(maxMs, Math.floor(baseMs * Math.pow(2, Math.max(0, attempt - 1))));

  // Add small jitter to spread bursts (±20%)
  const jitterFactor = 0.8 + Math.random() * 0.4;
  const expJittered = Math.max(0, Math.floor(exp * jitterFactor));

  if (hasExplicit) {
    // Add a small safety buffer to avoid retrying slightly too early
    const buffered = Math.max(0, Math.floor(explicitDelayMs + 50));
    return Math.min(maxMs, Math.max(expJittered, buffered));
  }

  // Fallback: at least 0.5s for the first retry
  return Math.min(maxMs, Math.max(500, expJittered));
}

/**
 * 从 429 错误中提取恢复时间戳(毫秒)
 * @param {Error} error - 错误对象
 * @returns {number|null} 恢复时间戳,如果无法解析返回 null
 */
function getUpstreamResetTimestamp(error) {
  const body = extractUpstreamErrorBody(error);
  const root = (body && typeof body === 'object') ? body : null;
  const inner = root?.error || root;
  const details = Array.isArray(inner?.details) ? inner.details : [];

  for (const d of details) {
    if (!d || typeof d !== 'object') continue;
    const meta = d.metadata && typeof d.metadata === 'object' ? d.metadata : null;
    const ts = meta?.quotaResetTimeStamp;
    if (typeof ts === 'string') {
      const t = Date.parse(ts);
      if (Number.isFinite(t)) {
        return t;
      }
    }
  }
  return null;
}

/**
 * 判断错误是否为可重试的临时性错误(429 或 503 容量不足)
 * @param {number} status - HTTP 状态码
 * @param {Error} error - 错误对象
 * @returns {boolean}
 */
function isRetryableError(status, error) {
  // 429 Rate Limit 总是可重试
  if (status === 429) return true;

  // 503 需要检查是否为容量不足错误
  if (status === 503) {
    const body = extractUpstreamErrorBody(error);
    const root = (body && typeof body === 'object') ? body : null;
    const inner = root?.error || root;
    const details = Array.isArray(inner?.details) ? inner.details : [];
    
    // 检查是否包含 MODEL_CAPACITY_EXHAUSTED
    for (const d of details) {
      if (d?.reason === 'MODEL_CAPACITY_EXHAUSTED') {
        return true;
      }
    }
  }

  return false;
}

/**
 * 带 429/503 重试的执行器
 * @param {Function} fn - 要执行的异步函数,接收 attempt 参数
 * @param {number} maxRetries - 最大重试次数
 * @param {Object} options - 可选参数
 * @param {string} options.loggerPrefix - 日志前缀
 * @param {Function} options.onAttempt - 每次尝试时的回调(用于记录请求次数)
 * @param {string} options.tokenId - Token ID(用于模型系列禁用)
 * @param {string} options.modelId - 模型 ID(用于模型系列禁用)
 * @param {Function} options.refreshQuota - 刷新额度的回调函数(当需要获取准确恢复时间时调用)
 * @returns {Promise<any>}
 */
export async function with429Retry(fn, maxRetries, options = {}, legacyOnAttempt = null) {
  // 兼容旧版调用方式:with429Retry(fn, maxRetries, loggerPrefix, onAttempt)
  let loggerPrefix = '';
  let onAttempt = null;
  let tokenId = null;
  let modelId = null;
  let refreshQuota = null;

  if (typeof options === 'string') {
    // 旧版调用方式
    loggerPrefix = options;
    onAttempt = legacyOnAttempt;
  } else if (typeof options === 'object' && options !== null) {
    loggerPrefix = options.loggerPrefix || '';
    onAttempt = options.onAttempt || null;
    tokenId = options.tokenId || null;
    modelId = options.modelId || null;
    refreshQuota = options.refreshQuota || null;
  }

  const retries = Number.isFinite(maxRetries) && maxRetries > 0 ? Math.floor(maxRetries) : 0;
  const cooldownThreshold = config.quota?.longCooldownThreshold || LONG_COOLDOWN_THRESHOLD;
  let attempt = 0;

  // 首次执行 + 最多 retries 次重试
  while (true) {
    try {
      // 每次尝试时调用回调(用于记录请求次数)
      if (typeof onAttempt === 'function') {
        onAttempt(attempt);
      }
      return await fn(attempt);
    } catch (error) {
      // 兼容多种错误格式:error.status, error.statusCode, error.response?.status
      const status = Number(error.status || error.statusCode || error.response?.status);

      if (isRetryableError(status, error)) {
        const explicitDelayMs = getUpstreamRetryDelayMs(error);
        const upstreamResetTimestamp = getUpstreamResetTimestamp(error);
        const errorType = status === 503 ? '503 (容量不足)' : '429';

        // 检查是否是长时间冷却(额度耗尽)- 仅 429 触发模型系列禁用
        if (status === 429 && explicitDelayMs !== null && explicitDelayMs >= cooldownThreshold && tokenId && modelId) {
          // 先检查是否已经被其他并发请求禁用了,避免重复处理
          if (!tokenCooldownManager.isAvailable(tokenId, modelId)) {
            // 已经在冷却中,直接抛出错误,不重复处理
            throw error;
          }

          // 恢复时间超过阈值,触发模型系列禁用
          // 优先使用上游返回的动态限流时间(更准确反映当前限流状态)
          let finalResetTimestamp = upstreamResetTimestamp;

          // 如果上游没有返回时间戳,使用延迟时长计算
          if (!finalResetTimestamp && explicitDelayMs !== null) {
            finalResetTimestamp = Date.now() + explicitDelayMs;
          }

          // 如果上游数据都没有,才尝试从 quotas.json 获取(作为兜底)
          if (!finalResetTimestamp && typeof refreshQuota === 'function') {
            logger.info(`${loggerPrefix}上游未返回恢复时间,尝试从额度数据获取...`);
            try {
              await refreshQuota();
              const { resetTime: quotaResetTime } = quotaManager.getModelGroupResetTime(tokenId, modelId);
              if (quotaResetTime) {
                finalResetTimestamp = quotaResetTime;
              }
            } catch (e) {
              logger.warn(`${loggerPrefix}获取额度数据失败: ${e.message}`);
            }
          }

          if (finalResetTimestamp && finalResetTimestamp > Date.now()) {
            const groupKey = getGroupKey(modelId);
            const resetDate = new Date(finalResetTimestamp);
            const delayMinutes = Math.round((finalResetTimestamp - Date.now()) / 1000 / 60);
          logger.warn(
            `${loggerPrefix}收到 ${errorType},恢复时间 ${delayMinutes} 分钟后,` +
              `超过阈值(${Math.round(cooldownThreshold / 1000 / 60)}分钟),` +
              `禁用 ${groupKey} 系列直到 ${resetDate.toLocaleString('zh-CN', { timeZone: 'Asia/Shanghai' })}`
            );
            tokenCooldownManager.setCooldown(tokenId, modelId, finalResetTimestamp);
            // 不重试,直接抛出错误
            throw error;
          }
        }

        // 短时间等待,正常重试
        if (attempt < retries) {
          const nextAttempt = attempt + 1;
          const waitMs = computeBackoffMs(nextAttempt, explicitDelayMs);
          logger.warn(
            `${loggerPrefix}收到 ${errorType},等待 ${waitMs}ms 后进行第 ${nextAttempt} 次重试(共 ${retries} 次)` +
            (explicitDelayMs !== null ? `(上游提示≈${explicitDelayMs}ms)` : '')
          );
          await sleep(waitMs);
          attempt = nextAttempt;
          continue;
        }
      }
      throw error;
    }
  }
};