6k6g-v3 / src /worker /queue.ts
fanyubo
diag: dequeue 时保存任务原始数据到 Redis 诊断 key
17a830e
Raw
History Blame Contribute Delete
10 kB
/**
* 任务轮询框架 - 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),
};
}