/** * 任务轮询框架 - 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 { 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 { 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 { 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 { 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), }; }