149 lines
3.1 KiB
JavaScript
149 lines
3.1 KiB
JavaScript
/**
|
|
* In-process subagent task registry (same loaded model; no second loadModel).
|
|
*/
|
|
|
|
const tasks = new Map();
|
|
const waiters = [];
|
|
const MAX_CONCURRENT = 2;
|
|
|
|
function makeId() {
|
|
return 'task_' + Date.now().toString(36) + '_' + Math.random().toString(36).slice(2, 8);
|
|
}
|
|
|
|
function create(label, extra) {
|
|
extra = extra || {};
|
|
if (runningCount() >= MAX_CONCURRENT) {
|
|
throw new Error('too many concurrent subagents (max ' + MAX_CONCURRENT + ')');
|
|
}
|
|
const id = makeId();
|
|
const kind = extra.subagentType === 'general' || extra.subagent_type === 'general' ? 'general' : 'explore';
|
|
tasks.set(id, {
|
|
id,
|
|
label: label || 'subagent',
|
|
status: 'running',
|
|
summary: null,
|
|
messages: [],
|
|
history: Array.isArray(extra.history) ? extra.history : [],
|
|
subagentType: kind,
|
|
createdAt: Date.now(),
|
|
});
|
|
return id;
|
|
}
|
|
|
|
function runningCount() {
|
|
let n = 0;
|
|
for (const t of tasks.values()) {
|
|
if (t.status === 'running') n += 1;
|
|
}
|
|
return n;
|
|
}
|
|
|
|
function setHistory(id, history) {
|
|
const t = tasks.get(id);
|
|
if (t) {
|
|
t.history = history || [];
|
|
t.updatedAt = Date.now();
|
|
}
|
|
return t || null;
|
|
}
|
|
|
|
function flushWaiters() {
|
|
for (let i = waiters.length - 1; i >= 0; i--) {
|
|
if (waiters[i]()) waiters.splice(i, 1);
|
|
}
|
|
}
|
|
|
|
function finish(id, summary) {
|
|
const t = tasks.get(id);
|
|
if (t) {
|
|
t.status = 'done';
|
|
t.summary = summary;
|
|
t.updatedAt = Date.now();
|
|
}
|
|
flushWaiters();
|
|
return t || null;
|
|
}
|
|
|
|
function fail(id, err) {
|
|
const t = tasks.get(id);
|
|
if (t) {
|
|
t.status = 'error';
|
|
t.summary = String(err || 'error');
|
|
t.updatedAt = Date.now();
|
|
}
|
|
flushWaiters();
|
|
return t || null;
|
|
}
|
|
|
|
function get(id) {
|
|
return tasks.get(id) || null;
|
|
}
|
|
|
|
function list() {
|
|
return Array.from(tasks.values()).map((t) => ({
|
|
id: t.id,
|
|
label: t.label,
|
|
status: t.status,
|
|
summary: t.summary,
|
|
subagentType: t.subagentType || 'explore',
|
|
}));
|
|
}
|
|
|
|
function waitAll(opts) {
|
|
opts = opts || {};
|
|
const timeoutMs = opts.timeoutMs > 0 ? opts.timeoutMs : 120000;
|
|
const running = () => list().filter((t) => t.status === 'running');
|
|
if (!running().length) return Promise.resolve({ running: 0, tasks: list() });
|
|
return new Promise((resolve) => {
|
|
const timer = setTimeout(() => {
|
|
resolve({ running: running().length, tasks: list(), timedOut: true });
|
|
}, timeoutMs);
|
|
waiters.push(() => {
|
|
if (!running().length) {
|
|
clearTimeout(timer);
|
|
resolve({ running: 0, tasks: list() });
|
|
return true;
|
|
}
|
|
return false;
|
|
});
|
|
});
|
|
}
|
|
|
|
function kill(id) {
|
|
const t = tasks.get(id);
|
|
if (t && t.status === 'running') {
|
|
t.status = 'killed';
|
|
t.updatedAt = Date.now();
|
|
}
|
|
flushWaiters();
|
|
return t || null;
|
|
}
|
|
|
|
function appendMessage(id, text) {
|
|
const t = tasks.get(id);
|
|
if (!t) throw new Error('unknown task: ' + id);
|
|
t.messages.push(String(text || ''));
|
|
t.updatedAt = Date.now();
|
|
return t;
|
|
}
|
|
|
|
function reset() {
|
|
tasks.clear();
|
|
waiters.length = 0;
|
|
}
|
|
|
|
module.exports = {
|
|
MAX_CONCURRENT,
|
|
create,
|
|
finish,
|
|
fail,
|
|
get,
|
|
list,
|
|
waitAll,
|
|
kill,
|
|
appendMessage,
|
|
setHistory,
|
|
runningCount,
|
|
reset,
|
|
};
|