Spaces:
Sleeping
Sleeping
| /** | |
| * 任务轮询框架 - Upstash REST 适配 | |
| * 从 Redis 队列获取慢任务(非阻塞轮询) | |
| * ★★★ 并发处理池:支持同时处理多个任务 ★★★ | |
| */ | |
| import { getRedis, RedisHelper } from '../lib/redis'; | |
| import { setWorkerRunning, setWorkerIdle, addActiveTask, removeActiveTask, WorkerStatus, getWorkerStatus } from './status'; | |
| import { processTask } from './processor'; | |
| import { setCurrentTaskForWatchdog, clearCurrentTaskForWatchdog, removeWatchedTask } from './index'; | |
| const QUEUE_KEYS = { | |
| boosted: 'slow_task:queue:boosted', | |
| normal: 'slow_task:queue:normal', | |
| results: 'slow_task:results', | |
| }; | |
| const POLL_INTERVAL = 3000; // 3秒轮询间隔 | |
| const MAX_CONCURRENT_TASKS = 3; // ★★★ 最大并发任务数 ★★★ | |
| // ★★★ 并发计数器 ★★★ | |
| let activeTasksCount = 0; | |
| // ★ 任务 2: 增强日志 - 轮询计数器 | |
| let pollCount = 0; | |
| let lastLogTime = Date.now(); | |
| /** | |
| * 获取当前活跃任务数 | |
| */ | |
| export function getActiveTasksCount(): number { | |
| return activeTasksCount; | |
| } | |
| /** | |
| * 获取最大并发数 | |
| */ | |
| export function getMaxConcurrentTasks(): number { | |
| return MAX_CONCURRENT_TASKS; | |
| } | |
| /** | |
| * 轮询任务队列(非阻塞) | |
| * Upstash REST 不支持 brpop,改用 rpop + 定时轮询 | |
| * ★★★ 并发模式:取出任务后立即异步处理,不阻塞轮询循环 ★★★ | |
| */ | |
| export async function pollTaskQueue(): Promise<void> { | |
| const redis = getRedis(); | |
| pollCount++; | |
| // ★ 每100次轮询打印一次状态(约每5分钟) | |
| if (pollCount % 100 === 0) { | |
| console.log(`[Queue] Poll #${pollCount}, active: ${activeTasksCount}/${MAX_CONCURRENT_TASKS}, elapsed: ${Math.floor((Date.now() - lastLogTime) / 1000)}s`); | |
| lastLogTime = Date.now(); | |
| } | |
| if (!redis) { | |
| console.warn('[Queue] ❌ Redis not available'); | |
| return; | |
| } | |
| // ★★★ 并发控制:如果已满,跳过本轮 ★★★ | |
| if (activeTasksCount >= MAX_CONCURRENT_TASKS) { | |
| return; | |
| } | |
| // 优先取加速队列,其次普通队列 | |
| let task: any = null; | |
| let queueName: string = ''; | |
| // 先检查加速队列 | |
| task = await RedisHelper.dequeue(QUEUE_KEYS.boosted); | |
| if (task) { | |
| queueName = QUEUE_KEYS.boosted; | |
| } else { | |
| // 再检查普通队列 | |
| task = await RedisHelper.dequeue(QUEUE_KEYS.normal); | |
| if (task) { | |
| queueName = QUEUE_KEYS.normal; | |
| } | |
| } | |
| if (!task) { | |
| // 无任务,等待下次轮询 | |
| return; | |
| } | |
| // ★ 诊断:保存任务原始数据到诊断 key(读取后可删除) | |
| try { | |
| const redis2 = getRedis(); | |
| if (redis2) { | |
| const diagData = { | |
| taskId: task.taskId, | |
| hasChatContext: !!task.chatContext, | |
| chatContextLen: task.chatContext?.length || 0, | |
| chatContext0Role: task.chatContext?.[0]?.role || 'N/A', | |
| chatContext0ContentLen: task.chatContext?.[0]?.content?.length || 0, | |
| chatContext0ContentPreview: (task.chatContext?.[0]?.content || '').slice(0, 300), | |
| promptLen: task.prompt?.length || 0, | |
| promptPreview: (task.prompt || '').slice(0, 300), | |
| hasDraftContent: !!task.draftContent, | |
| draftContentLen: task.draftContent?.length || 0, | |
| title: task.title, | |
| taskType: task.taskType, | |
| skipPhase1: task.skipPhase1, | |
| timestamp: Date.now(), | |
| }; | |
| await redis2.set('slow_task:diag:last_dequeued', JSON.stringify(diagData), 'EX', 3600); | |
| } | |
| } catch (diagErr) { | |
| // 忽略诊断错误 | |
| } | |
| console.log(`[Queue] Got task from: ${queueName}, active slots: ${activeTasksCount}/${MAX_CONCURRENT_TASKS}`); | |
| // ★★★ 并发控制:立即递增计数器 ★★★ | |
| activeTasksCount++; | |
| addActiveTask(task.taskId); | |
| // ★★★ 异步处理任务(不阻塞轮询循环)★★★ | |
| processTaskAsync(task, queueName).catch(() => { | |
| // 兜底:确保计数器不会泄漏 | |
| // processTaskAsync 内部已有完整错误处理,这里只是安全网 | |
| }); | |
| } | |
| /** | |
| * 异步处理单个任务(独立的错误边界) | |
| */ | |
| async function processTaskAsync(task: any, queueName: string): Promise<void> { | |
| try { | |
| console.log(`[Queue] ▶ Starting task: ${task.taskId} (type: ${task.taskType}), active: ${activeTasksCount}/${MAX_CONCURRENT_TASKS}`); | |
| // 设置看门狗监控任务 | |
| setCurrentTaskForWatchdog(task.taskId, task.taskType); | |
| await setWorkerRunning(task.taskId); | |
| // 处理任务 | |
| await processTask(task); | |
| console.log(`[Queue] ✅ Task completed: ${task.taskId}`); | |
| } catch (err: any) { | |
| console.error(`[Queue] ❌ Task failed: ${task.taskId} — ${err.message}`); | |
| console.error('[Queue] Error stack:', err.stack); | |
| // 将任务标记为失败 | |
| try { | |
| const redis = getRedis(); | |
| if (redis && task) { | |
| const existing = await redis.hget(QUEUE_KEYS.results, task.taskId); | |
| if (existing) { | |
| const state = typeof existing === 'string' ? JSON.parse(existing) : existing; | |
| const updated = { | |
| ...state, | |
| status: 'failed', | |
| errorMessage: err.message || 'Unknown error', | |
| errorStack: err.stack?.split('\n').slice(0, 3).join('\n'), | |
| updatedAt: Date.now(), | |
| }; | |
| await redis.hset(QUEUE_KEYS.results, { [task.taskId]: JSON.stringify(updated) }); | |
| console.log('[Queue] Task marked as failed:', task.taskId); | |
| } else { | |
| // 如果没有状态记录,创建一个 | |
| await redis.hset(QUEUE_KEYS.results, { | |
| [task.taskId]: JSON.stringify({ | |
| taskId: task.taskId, | |
| status: 'failed', | |
| errorMessage: err.message || 'Unknown error', | |
| updatedAt: Date.now(), | |
| }) | |
| }); | |
| console.log('[Queue] Created failed status for:', task.taskId); | |
| } | |
| } | |
| } catch (updateErr: any) { | |
| console.error('[Queue] Failed to update task status:', updateErr.message); | |
| } | |
| } finally { | |
| // ★★★ 无论成功失败,都递减计数器 ★★★ | |
| activeTasksCount--; | |
| removeActiveTask(task.taskId); | |
| removeWatchedTask(task.taskId); | |
| await setWorkerIdle(); | |
| console.log(`[Queue] ▶ Slot freed: ${activeTasksCount}/${MAX_CONCURRENT_TASKS} active`); | |
| } | |
| } | |
| /** | |
| * 启动轮询循环 | |
| * ★★★ 并发模式:轮询不阻塞,任务异步执行 ★★★ | |
| */ | |
| // ─── Zombie Recovery(OPT-07)─── | |
| const ZOMBIE_RECOVERY_THRESHOLD_MS = 15 * 60 * 1000; // 15 分钟 | |
| /** | |
| * 扫描 stuck 任务并重新入队 | |
| * Worker 启动时恢复因崩溃/重启而丢失的任务 | |
| */ | |
| async function recoverZombieTasks(): Promise<number> { | |
| const redis = getRedis(); | |
| if (!redis) return 0; | |
| try { | |
| const allResults = await redis.hgetall(QUEUE_KEYS.results); | |
| if (!allResults) return 0; | |
| const now = Date.now(); | |
| let recovered = 0; | |
| for (const [taskId, rawState] of Object.entries(allResults)) { | |
| const state = typeof rawState === 'string' ? JSON.parse(rawState) : rawState; | |
| if (!state) continue; | |
| if (state.status === 'completed' || state.status === 'failed') continue; | |
| const lastActive = state.updatedAt || state.startedAt || state.createdAt || 0; | |
| if (now - lastActive < ZOMBIE_RECOVERY_THRESHOLD_MS) continue; | |
| console.log(`[Queue] 🧟 发现 zombie 任务: ${taskId}, 停滞 ${Math.floor((now - lastActive) / 1000)}秒`); | |
| const recoveredState = { | |
| ...state, | |
| status: 'queued', | |
| progress: 0, | |
| currentStep: 0, | |
| detail: 'Zombie recovery: 任务已自动恢复', | |
| updatedAt: now, | |
| }; | |
| const queueKey = state.boosted ? QUEUE_KEYS.boosted : QUEUE_KEYS.normal; | |
| await redis.lpush(queueKey, JSON.stringify(recoveredState)); | |
| await redis.hset(QUEUE_KEYS.results, { [taskId]: JSON.stringify(recoveredState) }); | |
| recovered++; | |
| console.log(`[Queue] ✅ Zombie 任务已恢复: ${taskId}`); | |
| } | |
| return recovered; | |
| } catch (error) { | |
| console.warn('[Queue] Zombie recovery failed:', error); | |
| return 0; | |
| } | |
| } | |
| export async function startPollingLoop(): Promise<void> { | |
| console.log(`[Queue] Starting concurrent polling loop, interval: ${POLL_INTERVAL}ms, max concurrent: ${MAX_CONCURRENT_TASKS}`); | |
| // ★ OPT-07: Worker 启动时执行 zombie recovery | |
| try { | |
| const recovered = await recoverZombieTasks(); | |
| if (recovered > 0) { | |
| console.log(`[Queue] 🔄 Zombie recovery: 恢复了 ${recovered} 个任务`); | |
| } | |
| } catch (err) { | |
| console.warn('[Queue] Zombie recovery skipped:', err); | |
| } | |
| // ★★★ 轮询心跳:确保循环在运行 ★★★ | |
| let loopCount = 0; | |
| const LOOP_HEARTBEAT_KEY = 'slow_task:poll_heartbeat'; | |
| while (true) { | |
| loopCount++; | |
| // 每10次轮询更新心跳(约30秒) | |
| if (loopCount % 10 === 0) { | |
| const redis = getRedis(); | |
| if (redis) { | |
| try { | |
| const status = getWorkerStatus(); | |
| await redis.set(LOOP_HEARTBEAT_KEY, JSON.stringify({ | |
| loopCount, | |
| timestamp: Date.now(), | |
| activeTasks: activeTasksCount, | |
| maxConcurrent: MAX_CONCURRENT_TASKS, | |
| processedCount: status.processedCount | |
| }), 'EX', 300); | |
| console.log(`[Queue] 💓 Poll heartbeat: loop #${loopCount}, active: ${activeTasksCount}/${MAX_CONCURRENT_TASKS}`); | |
| } catch (err) { | |
| console.warn('[Queue] Heartbeat write failed'); | |
| } | |
| } | |
| } | |
| // ★★★ 关键变更:不 await,让任务异步执行 ★★★ | |
| // pollTaskQueue 内部有并发控制,满了就跳过 | |
| await pollTaskQueue(); | |
| // 等待下次轮询 | |
| await new Promise(resolve => setTimeout(resolve, POLL_INTERVAL)); | |
| } | |
| } | |
| /** | |
| * 获取队列状态 | |
| */ | |
| export async function getQueueStatus(): Promise<{ | |
| boosted: number; | |
| normal: number; | |
| }> { | |
| const redis = getRedis(); | |
| if (!redis) { | |
| return { boosted: 0, normal: 0 }; | |
| } | |
| return { | |
| boosted: await RedisHelper.queueLength(QUEUE_KEYS.boosted), | |
| normal: await RedisHelper.queueLength(QUEUE_KEYS.normal), | |
| }; | |
| } |