BackgroundQueue 为每次失败的尝试留下一个待定的超时计时器,阻止进程退出
const logger = { trace() {}, debug() {}, error() {} };
export class BackgroundQueue {
tasks = [];
activeTasks = new Set();
constructor(options = {}) {
this.options = {
maxConcurrency: options.maxConcurrency ?? 3,
defaultTimeout: options.defaultTimeout ?? 10000,
defaultRetries: options.defaultRetries ?? 2,
};
}
enqueue(task) {
task.timeout = task.timeout ?? this.options.defaultTimeout;
task.retries = task.retries ?? this.options.defaultRetries;
this.tasks.push(task);
logger.trace(Enqueued task ${task.id});
setTimeout(() => this.processNext(), 0);
}
processNext() {
while (this.tasks.length > 0 && this.activeTasks.size < this.options.maxConcurrency) {
const task = this.tasks.shift();
if (!task) break;
const taskPromise = this.executeTask(task);
this.activeTasks.add(taskPromise);
taskPromise.finally(() => {
this.activeTasks.delete(taskPromise);
setTimeout(() => this.processNext(), 0);
});
}
}
async executeTask(task) {
let lastError;
const maxAttempts = (task.retries ?? 0) + 1;
for (let attempt = 1; attempt <= maxAttempts; attempt++) {
try {
let timeoutId;
const timeoutPromise = new Promise((_, reject) => {
timeoutId = setTimeout(() => {
reject(new Error(Task ${task.id} timeout));
}, task.timeout);
});
const result = await Promise.race([task.operation(), timeoutPromise]);
if (timeoutId) clearTimeout(timeoutId);
logger.trace(Task ${task.id} completed (attempt ${attempt}/${maxAttempts});
return result;
} catch (error) {
lastError = error instanceof
内容来源: VoltAgent/voltagent