Files
gnome-jarvis/vendor/agent-harness/agent/loop.js
T
snxraven 9ad09b593c
Rolling release / release (push) Successful in 7m7s
Updates
2026-09-12 17:04:03 -04:00

1154 lines
42 KiB
JavaScript

/**
* Agentic turn loop: sample (QVAC) → tools → repeat.
*/
const engine = require('../lib/qvac.js');
const catalog = require('../lib/catalog.js');
const toolParse = require('../lib/tool-parse.js');
const sessions = require('./sessions.js');
const tools = require('./tools.js');
const sandbox = require('./sandbox.js');
const prompts = require('./prompts.js');
const compaction = require('./compaction.js');
const tasks = require('./tasks.js');
const customTools = require('./custom-tools.js');
const toolSet = require('./tool-set.js');
const planMode = require('./plan-mode.js');
const todos = require('./todos.js');
const stationarity = require('./stationarity.js');
const toolBudget = require('./tool-budget.js');
const goalMod = require('./goal.js');
const truncate = require('./truncate.js');
const sr = require('./search-replace.js');
const toolBatch = require('./tool-batch.js');
const gitSidecar = require('./git-sidecar.js');
const permStore = require('./perm-store.js');
const permRules = require('./perm-rules.js');
const memory = require('./memory.js');
const mcp = require('./mcp.js');
const policy = require('./policy.js');
const path = require('path');
const fs = require('fs');
customTools.setReserved(toolSet.ALWAYS_RESERVED.concat(toolSet.ALWAYS_BUILTIN_RESERVED));
const live = new Map();
const pendingPerms = new Map();
const pendingCustom = new Map();
const pendingAsks = new Map();
const pendingPlans = new Map();
const MAX_TURNS = 24;
const SUBAGENT_TURNS = 8;
const CUSTOM_TOOL_TIMEOUT_MS = 60000;
const ASK_TIMEOUT_MS = 10 * 60 * 1000;
const PLAN_TIMEOUT_MS = 10 * 60 * 1000;
const MAX_GOAL_NUDGES = 8;
function fsRead(cwd, rel) {
return fs.readFileSync(path.join(cwd, rel), 'utf8');
}
async function ensureModel(model) {
const loaded = engine.getLoaded();
const id = model || loaded.friendlyId || 'qwen3.5-4b';
const entry = catalog.findCatalogEntry(id);
const same =
loaded.modelId && (!model || loaded.friendlyId === model || loaded.constant === model);
if (same && !(entry && entry.vision && !loaded.vision)) return loaded;
return engine.load({ model: id, tools: true, device: 'auto' });
}
function emitUpdate(emit, sessionId, jobId, update) {
sessions.appendUpdate(sessionId, update);
emit('cap-chunk', Object.assign({ pack: 'agent', sessionId, jobId, kind: 'session_update' }, update));
}
function emitLive(emit, sessionId, jobId, update) {
emit('cap-chunk', Object.assign({ pack: 'agent', sessionId, jobId, kind: 'session_update' }, update));
}
function loadedCtxSize() {
const loaded = engine.getLoaded && engine.getLoaded();
return (loaded && loaded.ctxSize) || 8192;
}
function toolResultCap(budget) {
const ctx = loadedCtxSize();
const voice = !!(budget && budget.voice);
const max = voice ? 2000 : 8000;
const ratio = voice ? 0.1 : 0.35;
const min = voice ? 600 : 1200;
return Math.min(max, Math.max(min, Math.floor(ctx * ratio)));
}
function emitCompactDone(emit, session, jobId, toolDefs, ctxSize, beforeUsage, method) {
const afterUsage = compaction.usage(session.history, toolDefs, ctxSize);
emitUpdate(emit, session.id, jobId, {
type: 'compaction',
status: 'done',
method: method,
used: afterUsage.used,
limit: afterUsage.limit,
pct: afterUsage.pct,
before: beforeUsage && beforeUsage.used,
threshold: afterUsage.threshold,
});
emitLive(emit, session.id, jobId, Object.assign({ type: 'context' }, afterUsage));
return afterUsage;
}
function usageFromStats(stats, fallbackUsed, ctxSize) {
if (!stats || typeof stats !== 'object') return compaction.snapshot(fallbackUsed, ctxSize);
const prompt = Number(
stats.n_past != null
? stats.n_past
: stats.cacheTokens != null
? stats.cacheTokens
: stats.prompt_n != null
? Number(stats.prompt_n) + Number(stats.predicted_n || stats.n_predicted || 0)
: stats.tokens != null
? stats.tokens
: NaN
);
if (!Number.isFinite(prompt) || prompt <= 0) return compaction.snapshot(fallbackUsed, ctxSize);
return compaction.snapshot(prompt, ctxSize);
}
function withUsage(update, used, ctxSize) {
return Object.assign(update, compaction.snapshot(used, ctxSize));
}
function persistSession(session, tracker) {
if (tracker) session.planMode = planMode.snapshot(tracker);
if (session.goal) session.goal = goalMod.snapshot(session.goal);
sessions.saveSummary(session);
}
async function waitKeyed(map, key, timeoutMs, timeoutValue) {
const existing = map.get(key);
if (existing && existing.ready) {
map.delete(key);
return existing.ready;
}
return new Promise((resolve) => {
const timer = setTimeout(() => {
map.delete(key);
resolve(timeoutValue);
}, timeoutMs);
map.set(key, {
resolve(payload) {
clearTimeout(timer);
resolve(payload);
},
});
});
}
function resolveKeyed(map, key, payload) {
const rec = map.get(key);
if (rec && typeof rec.resolve === 'function') {
map.delete(key);
rec.resolve(payload);
return true;
}
map.set(key, { ready: payload });
return true;
}
async function waitPermission(jobId, payload) {
return new Promise((resolve) => {
pendingPerms.set(jobId + ':' + payload.toolCallId, resolve);
});
}
function resolvePermission(jobId, toolCallId, decision) {
const key = jobId + ':' + toolCallId;
const fn = pendingPerms.get(key);
if (fn) {
pendingPerms.delete(key);
fn(decision);
return true;
}
const keys = Array.from(pendingPerms.keys());
for (const k of keys) {
if (k === toolCallId || k.endsWith(':' + toolCallId)) {
const resolve = pendingPerms.get(k);
pendingPerms.delete(k);
if (resolve) resolve(decision);
return true;
}
}
return false;
}
function waitCustomResult(jobId, toolCallId) {
const key = jobId + ':' + toolCallId;
const existing = pendingCustom.get(key);
if (existing && existing.ready) {
pendingCustom.delete(key);
return Promise.resolve(existing.ready);
}
return new Promise((resolve) => {
const timer = setTimeout(() => {
pendingCustom.delete(key);
resolve({ error: 'custom tool timed out' });
}, CUSTOM_TOOL_TIMEOUT_MS);
pendingCustom.set(key, {
resolve(payload) {
clearTimeout(timer);
resolve(payload);
},
});
});
}
function resolveCustomResult(jobId, toolCallId, payload) {
return resolveKeyed(pendingCustom, jobId + ':' + toolCallId, payload || {});
}
function waitAsk(jobId, toolCallId) {
return waitKeyed(pendingAsks, jobId + ':' + toolCallId, ASK_TIMEOUT_MS, { error: 'ask_user timed out' });
}
function resolveAsk(jobId, toolCallId, choice) {
const payload = { choice };
if (jobId) {
const key = jobId + ':' + toolCallId;
if (pendingAsks.has(key)) return resolveKeyed(pendingAsks, key, payload);
}
const keys = Array.from(pendingAsks.keys());
for (const key of keys) {
if (key === toolCallId || key.endsWith(':' + toolCallId)) return resolveKeyed(pendingAsks, key, payload);
}
return resolveKeyed(pendingAsks, (jobId || '') + ':' + toolCallId, payload);
}
function waitPlanDecision(sessionId) {
return waitKeyed(pendingPlans, sessionId, PLAN_TIMEOUT_MS, { decision: 'timeout' });
}
function resolvePlanDecision(sessionId, decision) {
return resolveKeyed(pendingPlans, sessionId, { decision: decision === 'approve' ? 'approve' : 'reject' });
}
function buildToolDefs(session, payload, tracker) {
const hostWorkspace = session.hostWorkspace !== false;
const defs = tools
.defs({
planMode: planMode.isActive(tracker),
webFetch: payload && payload.webFetch,
hostWorkspace,
builtinTools: session.builtinTools,
})
.concat(customTools.defs(session.id));
return catalog.filterToolsForModel(defs, session.model);
}
function refreshSystem(session, sys) {
if (!session.history) session.history = [];
if (session.history.length && session.history[0].role === 'system') {
session.history[0] = { role: 'system', content: sys };
} else {
session.history.unshift({ role: 'system', content: sys });
}
}
function pushHistory(session, msg) {
session.history.push(msg);
if (msg.role !== 'system') sessions.appendHistory(session.id, msg);
}
function applyPlanWrite(session, name, args) {
let text = sessions.readPlan(session.id) || '';
if (name === 'write_file') {
text = args.contents != null ? String(args.contents) : args.content != null ? String(args.content) : '';
} else {
const old = args.old_string || args.oldString || '';
const neu = args.new_string != null ? args.new_string : args.newString;
if (neu == null) throw new Error('new_string required');
const applied = sr.applySearchReplace(text, old, neu, !!(args.replace_all || args.replaceAll));
text = applied.text;
}
const file = sessions.writePlan(session.id, text);
if (session.hostWorkspace !== false && session.cwd) {
try {
const abs = path.join(session.cwd, 'plan.md');
fs.writeFileSync(abs, text);
} catch (_) {}
}
return { ok: true, path: file, bytes: text.length };
}
function appendCompactReminders(extra, session, tracker, budget) {
if (!tracker || !tracker.pendingCompactReminder) return;
const voice = !!(budget && budget.voice);
const reminders = [{ role: 'user', content: compaction.compactReminder({ voice }) }];
const continuation = compaction.autoContinue(session.history, { voice });
if (continuation) reminders.push(continuation);
for (const reminder of reminders) {
if (!extra.some((m) => m.content === reminder.content)) extra.push(reminder);
}
tracker.pendingCompactReminder = false;
}
function sidecarMessages(session, tracker, budget) {
const extra = [];
if (budget && budget.answerOnly) {
extra.push({ role: 'user', content: toolBudget.answerNowMessage() });
}
appendCompactReminders(extra, session, tracker, budget);
if (planMode.isActive(tracker)) {
extra.push({
role: 'user',
content:
'<system-reminder>\n' +
planMode.reminder(tracker, { hasContent: !!(sessions.readPlan(session.id) || '').trim() }) +
'\n</system-reminder>',
});
} else if (tracker && tracker.pendingExitReminder) {
extra.push({ role: 'user', content: '<system-reminder>\n' + planMode.exitReminder() + '\n</system-reminder>' });
tracker.pendingExitReminder = false;
}
if (session.plan && session.plan.length) {
extra.push({ role: 'user', content: todos.formatBlock(session.plan) });
}
const mcpNotes = mcp.handshakeReminders();
for (const note of mcpNotes) {
extra.push({ role: 'user', content: '<system-reminder>\n' + note + '\n</system-reminder>' });
}
return extra;
}
function endTurn(emit, session, jobId, tracker, extra) {
persistSession(session, tracker);
const payload = Object.assign({ type: 'end', reason: 'stop' }, extra || {});
payload.context = compaction.usage(session.history, null, loadedCtxSize());
emitUpdate(emit, session.id, jobId, payload);
return {
ok: payload.reason !== 'cancelled',
text: payload.text || '',
turns: payload.turns,
reason: payload.reason,
};
}
async function verifyGoal(session, tracker, lastText) {
const g = session.goal;
if (!g || g.verify === false) {
g.status = 'complete';
return { achieved: true, gaps: [] };
}
if ((g.verifierRuns || 0) >= goalMod.MAX_VERIFIER_RUNS) {
return { achieved: false, gaps: ['verification budget exhausted'] };
}
g.status = 'verifying';
g.verifierRuns = (g.verifierRuns || 0) + 1;
const evidence = [
lastText || '',
todos.formatBlock(session.plan),
(sessions.readPlan(session.id) || '').slice(0, 8000),
g.notes || '',
]
.filter(Boolean)
.join('\n\n');
try {
const result = await engine.complete({
history: [
{ role: 'system', content: 'Reply with JSON only.' },
{ role: 'user', content: goalMod.verifierPrompt(g, evidence) },
],
tools: [],
desktopVision: false,
});
const verdict = goalMod.parseVerifier(result && result.text);
g.gaps = verdict.gaps || [];
if (verdict.achieved) g.status = 'complete';
else g.status = 'executing';
persistSession(session, tracker);
return verdict;
} catch (err) {
g.status = 'executing';
g.gaps = ['verifier failed: ' + err.message];
persistSession(session, tracker);
return { achieved: false, gaps: g.gaps };
}
}
async function runTurn(ctx) {
const { session, userText, emit, jobId, payload } = ctx;
const origin = session.origin || (payload && payload._origin) || '';
const hostWorkspace = session.hostWorkspace !== false;
const cwd = hostWorkspace ? session.cwd || sandbox.defaultCwd(origin) : session.cwd || 'page';
const mode = sandbox.permissionMode(payload);
const tracker = planMode.create(session.planMode);
if (payload && (payload.planMode === true || payload.planMode === 'active' || payload.planMode === 'plan')) {
if (tracker.state === 'inactive') {
planMode.enterPending(tracker);
planMode.activate(tracker);
}
}
ctx.planTracker = tracker;
ctx.planMode = planMode.isActive(tracker);
if (payload && payload.goal) {
session.goal = goalMod.create(payload.goal, { verify: payload.verify !== false });
}
if (payload && payload.verify === false && session.goal) session.goal.verify = false;
await ensureModel(session.model);
const voice = (payload && payload.voice === true) || origin === 'jarvis-qvac';
let extraSys = payload && payload.system;
if (session.goal && goalMod.isActive(session.goal)) {
extraSys = [extraSys, goalMod.plannerAddendum(session.goal)].filter(Boolean).join('\n\n');
}
const sysBase = prompts.assemble({
cwd: hostWorkspace ? cwd : session.workspace || cwd,
hostWorkspace,
extra: extraSys,
personality: voice ? 'voice' : undefined,
fsRead:
hostWorkspace
? (c, r) => {
try {
return fsRead(c, r);
} catch (_) {
return '';
}
}
: null,
});
const sys = catalog.isCompactToolModel(session.model)
? sysBase + '\n\n' + toolParse.FORMAT_REMINDER
: sysBase;
const sidecars = [];
if (hostWorkspace) {
if (!voice) {
const gitText = await gitSidecar.gitStatusSb(cwd, { hostWorkspace: true, run: tools.runShell });
if (gitText) sidecars.push(gitText);
}
const memText = memory.injectBlock(origin, userText || '');
if (memText) sidecars.push(memText);
}
refreshSystem(session, sidecars.length ? sys + '\n\n' + sidecars.join('\n\n') : sys);
if (userText || (payload && payload.images && payload.images.length)) {
let userMsg = { role: 'user', content: userText || '' };
if (payload && payload.images && payload.images.length) {
userMsg = engine.prepareVisionHistory(
[{ role: 'user', content: userText || '', images: payload.images }],
{ origin, cwd }
)[0];
userMsg = {
role: 'user',
content: userMsg.content,
attachments: userMsg.attachments,
};
}
pushHistory(session, userMsg);
}
persistSession(session, tracker);
if (session.goal) {
emitUpdate(emit, session.id, jobId, { type: 'goal_update', goal: goalMod.snapshot(session.goal) });
}
const cancelled = () => live.get(session.id) && live.get(session.id).cancelled;
const stuck = stationarity.create();
const budget = toolBudget.fromPayload(payload, origin);
let lastText = '';
let goalNudges = 0;
let toolNudges = 0;
async function runOneTool(item, turn) {
const name = item.name;
const args = item.args;
const toolCallId = item.toolCallId;
if (cancelled()) return { reason: 'cancelled', turns: turn + 1 };
let out;
try {
if (
planMode.isActive(tracker) &&
(name === 'write_file' || name === 'search_replace') &&
planMode.isPlanFilePath(args.path || args.file, tracker.planPath)
) {
out = applyPlanWrite(session, name, args);
emitUpdate(emit, session.id, jobId, {
type: 'plan_update',
path: tracker.planPath,
text: sessions.readPlan(session.id),
});
} else if (customTools.has(session.id, name)) {
const handler = customTools.getHandler(session.id, name);
if (typeof handler === 'function') {
out = await handler(args, { sessionId: session.id, toolCallId, name });
} else {
emitUpdate(emit, session.id, jobId, {
type: 'tool_request',
toolCallId,
name,
args,
});
const answered = await waitCustomResult(jobId, toolCallId);
if (answered && answered.error) throw new Error(answered.error);
out = answered && answered.result != null ? answered.result : answered;
}
} else if (name === 'task') {
out = await runSubagent(ctx, args);
} else if (name === 'send_subagent_message') {
const rec = tasks.appendMessage(args.task_id || args.taskId, args.message);
out = await continueSubagent(ctx, rec, args.message);
} else {
out = await tools.execute(
{
origin,
cwd,
session,
planMode: planMode.isActive(tracker),
planTracker: tracker,
hostWorkspace,
},
name,
args
);
if (out && out.type === 'enter_plan_mode') {
persistSession(session, tracker);
}
if (out && out.type === 'exit_plan_mode') {
planMode.requestExit(tracker);
persistSession(session, tracker);
emitUpdate(emit, session.id, jobId, {
type: 'plan_approval',
toolCallId,
plan: sessions.readPlan(session.id),
});
const ans = await waitPlanDecision(session.id);
if (cancelled()) return { reason: 'cancelled', turns: turn + 1 };
if (ans && ans.decision === 'approve') {
planMode.approveExit(tracker);
tracker.pendingExitReminder = true;
persistSession(session, tracker);
out = { planMode: false, approved: true, message: planMode.exitReminder() };
} else {
planMode.rejectExit(tracker);
persistSession(session, tracker);
out = {
planMode: true,
approved: false,
message: 'Plan rejected. Stay in plan mode and revise plan.md.',
};
}
}
if (out && out.type === 'ask_user') {
emitUpdate(emit, session.id, jobId, {
type: 'ask_user',
toolCallId,
question: out.question,
options: out.options,
});
const ans = await waitAsk(jobId, toolCallId);
if (cancelled()) return { reason: 'cancelled', turns: turn + 1 };
if (ans && ans.error) out = { error: ans.error };
else out = { type: 'ask_user', question: out.question, choice: ans && ans.choice };
}
if (out && out.type === 'goal_blocked') {
session.goal.status = 'blocked';
emitUpdate(emit, session.id, jobId, { type: 'goal_update', goal: goalMod.snapshot(session.goal) });
const rendered = truncate.renderToolResult(out, toolResultCap(budget));
pushHistory(session, { role: 'tool', name, content: rendered, tool_call_id: toolCallId });
emitUpdate(emit, session.id, jobId, { type: 'tool_result', toolCallId, name, result: rendered.slice(0, 4000) });
return { reason: 'goal_blocked', text: out.blocked_reason || lastText, turns: turn + 1 };
}
if (out && out.type === 'goal_completed') {
emitUpdate(emit, session.id, jobId, { type: 'goal_update', goal: goalMod.snapshot(session.goal) });
const skipVerify = payload && payload.verify === false;
const verdict = skipVerify ? { achieved: true, gaps: [] } : await verifyGoal(session, tracker, lastText);
if (verdict.achieved) {
const rendered = truncate.renderToolResult({ ok: true, achieved: true }, toolResultCap(budget));
pushHistory(session, { role: 'tool', name, content: rendered, tool_call_id: toolCallId });
emitUpdate(emit, session.id, jobId, { type: 'tool_result', toolCallId, name, result: rendered });
emitUpdate(emit, session.id, jobId, { type: 'goal_update', goal: goalMod.snapshot(session.goal) });
return { reason: 'goal_complete', text: lastText, turns: turn + 1 };
}
out = {
ok: false,
achieved: false,
gaps: verdict.gaps,
message: 'Verifier rejected completion. Fix the gaps and continue.',
};
emitUpdate(emit, session.id, jobId, { type: 'goal_update', goal: goalMod.snapshot(session.goal) });
}
}
} catch (err) {
out = { error: err.message };
}
const rendered = truncate.renderToolResult(out, toolResultCap(budget));
pushHistory(session, { role: 'tool', name, content: rendered, tool_call_id: toolCallId });
emitUpdate(
emit,
session.id,
jobId,
withUsage(
{ type: 'tool_result', toolCallId, name, result: rendered.slice(0, 4000) },
compaction.usage(session.history, null, loadedCtxSize()).used,
loadedCtxSize()
)
);
return null;
}
try {
for (let turn = 0; turn < budget.maxTurns; turn++) {
if (cancelled()) return endTurn(emit, session, jobId, tracker, { reason: 'cancelled', turns: turn });
if (turn === budget.maxTurns - 1) toolBudget.forceAnswer(budget);
const toolDefs = budget.answerOnly ? [] : buildToolDefs(session, payload, tracker);
ctx.planMode = planMode.isActive(tracker);
const ctxSize = loadedCtxSize();
// Generate one-shot reminders once; usage estimation must not consume
// them before inference. Reserve their space during compaction as well.
const turnSidecars = sidecarMessages(session, tracker, budget);
const sidecarTokens = compaction.conversationTokens(turnSidecars) + 256;
const beforeUsage = compaction.usage(session.history, toolDefs, ctxSize);
emitLive(emit, session.id, jobId, Object.assign({ type: 'context' }, beforeUsage));
if (compaction.shouldCompact(session.history, toolDefs, ctxSize, sidecarTokens)) {
emitUpdate(emit, session.id, jobId, {
type: 'compaction',
status: 'start',
method: 'llm',
used: beforeUsage.used,
limit: beforeUsage.limit,
pct: beforeUsage.pct,
threshold: beforeUsage.threshold,
});
const compactOpts = {
budgetTokens: compaction.historyBudget(ctxSize, toolDefs, 0, sidecarTokens),
tools: toolDefs,
voice: !!budget.voice,
};
session.history = await compaction.compactWithLlm(session.history, Object.assign({}, compactOpts, {
complete: (opts) => engine.complete(Object.assign({}, opts, {
desktopVision: false,
timeoutMs: budget.completeTimeoutMs,
idleMs: budget.completeIdleMs,
})),
}));
sessions.replaceHistory(session.id, session.history);
tracker.pendingCompactReminder = true;
if (hostWorkspace) {
const memText = memory.injectBlock(origin, userText || '');
if (memText) {
const head = session.history[0] && session.history[0].role === 'system' ? session.history[0].content : sys;
if (String(head).indexOf('[memory]') < 0) refreshSystem(session, head + '\n\n' + memText);
}
}
const afterUsage = compaction.usage(session.history, toolDefs, ctxSize);
emitUpdate(emit, session.id, jobId, {
type: 'compaction',
status: 'done',
method: 'llm',
used: afterUsage.used,
limit: afterUsage.limit,
pct: afterUsage.pct,
before: beforeUsage.used,
threshold: afterUsage.threshold,
});
emitLive(emit, session.id, jobId, Object.assign({ type: 'context' }, afterUsage));
}
if (cancelled()) return endTurn(emit, session, jobId, tracker, { reason: 'cancelled', turns: turn });
appendCompactReminders(turnSidecars, session, tracker, budget);
emitUpdate(emit, session.id, jobId, { type: 'turn', turn });
let streamBase = compaction.usage(session.history.concat(turnSidecars), toolDefs, ctxSize);
emitLive(emit, session.id, jobId, Object.assign({ type: 'context' }, streamBase));
let streamChars = 0;
function liveUsed() {
return streamBase.used + Math.ceil(streamChars / compaction.CHAR_PER_TOKEN);
}
function emitStream(update, moreText) {
if (moreText) streamChars += String(moreText).length;
emitUpdate(emit, session.id, jobId, withUsage(update, liveUsed(), ctxSize));
}
let result;
for (let overflowTry = 0; overflowTry < 4; overflowTry++) {
const history = session.history.concat(turnSidecars);
streamBase = compaction.usage(history, toolDefs, ctxSize);
streamChars = 0;
try {
result = await engine.complete(
{
history,
tools: toolDefs,
toolDialect: catalog.toolDialectFor(session.model),
desktopVision: payload && payload.desktopVision === false ? false : undefined,
timeoutMs: budget.completeTimeoutMs,
idleMs: budget.completeIdleMs,
},
(ev) => {
if (ev.type === 'contentDelta') {
emitStream({ type: 'agent_message_chunk', text: ev.delta }, ev.delta);
} else if (ev.type === 'thinkingDelta') {
emitStream({ type: 'agent_thought_chunk', text: ev.delta }, ev.delta);
} else if (ev.type === 'toolCall') {
const call = ev.call || {};
let extra = String(call.name || '');
try {
extra +=
typeof call.arguments === 'string'
? call.arguments
: JSON.stringify(call.arguments || call.args || {});
} catch (_) {}
emitStream({ type: 'tool_call', call: call }, extra);
}
}
);
break;
} catch (err) {
if (!compaction.isOverflowError(err) || overflowTry === 3) throw err;
emitUpdate(emit, session.id, jobId, {
type: 'compaction',
status: 'start',
method: 'overflow',
used: streamBase.used,
limit: streamBase.limit,
pct: streamBase.pct,
threshold: streamBase.threshold,
});
session.history = compaction.compact(session.history, {
// The estimate can be lower than the model's tokenizer count.
// Every overflow retry must shrink even an apparently small history.
budgetTokens: Math.min(
compaction.historyBudget(ctxSize, toolDefs, overflowTry + 1, sidecarTokens),
Math.max(1, Math.floor(compaction.conversationTokens(session.history) * 0.85))
),
tools: toolDefs,
aggressive: true,
voice: !!budget.voice,
});
sessions.replaceHistory(session.id, session.history);
tracker.pendingCompactReminder = true;
appendCompactReminders(turnSidecars, session, tracker, budget);
emitCompactDone(emit, session, jobId, toolDefs, ctxSize, streamBase, 'overflow');
}
}
emitLive(
emit,
session.id,
jobId,
Object.assign({ type: 'context' }, usageFromStats(result && result.stats, liveUsed(), ctxSize))
);
let calls = (result && result.toolCalls) || [];
if (!calls.length && result) {
const recovered = toolParse.recover({
text: result.text,
thinking: result.thinking,
tools: toolDefs,
existing: calls,
});
calls = recovered.calls;
if (recovered.text != null) result.text = recovered.text;
} else if (result && result.text) {
result.text = toolParse.stripToolMarkup(result.text);
}
if (budget.answerOnly) calls = [];
if ((result && result.text) || calls.length) {
if (result.text) lastText = result.text;
const assistant = { role: 'assistant', content: result.text || '' };
if (calls.length) assistant.tool_calls = calls;
pushHistory(session, assistant);
}
if (!calls.length) {
const goalActive = goalMod.isActive(session.goal);
if (goalActive && goalNudges < MAX_GOAL_NUDGES && !budget.answerOnly) {
goalNudges += 1;
pushHistory(session, { role: 'user', content: goalMod.continuation(session.goal) });
continue;
}
if (toolNudges < 2 && toolBudget.shouldNudgeToolCall([result && result.text, result && result.thinking].filter(Boolean).join('\n'), budget)) {
toolNudges += 1;
pushHistory(session, { role: 'user', content: toolBudget.continueToolMessage() });
continue;
}
return endTurn(emit, session, jobId, tracker, {
type: 'end',
reason: 'stop',
text: (result && result.text) || lastText || toolBudget.lastToolText(session.history) || '',
turns: turn + 1,
});
}
let stuckNow = false;
const prepared = [];
for (const call of calls) {
if (cancelled()) return endTurn(emit, session, jobId, tracker, { reason: 'cancelled', turns: turn + 1 });
const name = call.name;
const args = typeof call.arguments === 'string' ? safeJson(call.arguments) : call.arguments || {};
const toolCallId = call.id || name + '_' + Date.now();
stationarity.observe(stuck, name, args);
if (stationarity.shouldStop(stuck)) {
stuckNow = true;
const msg = 'Repeating the same tool call; stopping.';
pushHistory(session, { role: 'tool', name, content: msg, tool_call_id: toolCallId });
emitUpdate(emit, session.id, jobId, { type: 'tool_result', toolCallId, name, result: msg });
break;
}
const gateErr = planMode.gateWrite(tracker, name, args);
if (gateErr) {
pushHistory(session, { role: 'tool', name, content: gateErr, tool_call_id: toolCallId });
emitUpdate(emit, session.id, jobId, { type: 'tool_result', toolCallId, name, result: gateErr });
continue;
}
if (name === 'run_terminal_cmd' && toolBudget.shouldSkipShell(budget)) {
const skipped = toolBudget.skipShellMessage();
pushHistory(session, { role: 'tool', name, content: skipped, tool_call_id: toolCallId });
emitUpdate(emit, session.id, jobId, { type: 'tool_result', toolCallId, name, result: skipped });
continue;
}
if (customTools.needsPermission(session.id, name, mode) || (sandbox.needsPermission(name, mode) && !((name === 'write_file' || name === 'search_replace') && (planMode.isPlanFilePath(args.path, tracker.planPath) || policy.isIdentityPath(args.path || args.file))))) {
let remembered = null;
try {
remembered = permStore.resolve(name, args);
} catch (_) {}
if (remembered === 'deny') {
const denied = 'permission denied for ' + name;
pushHistory(session, { role: 'tool', name, content: denied, tool_call_id: toolCallId });
emitUpdate(emit, session.id, jobId, { type: 'tool_result', toolCallId, name, result: denied });
continue;
}
if (remembered !== 'allow' && mode !== 'always-approve') {
emitUpdate(emit, session.id, jobId, {
type: 'permission',
toolCallId,
tool: name,
args,
pattern: permRules.patternFromArgs(name, args),
});
let decision = await waitPermission(jobId, { toolCallId });
if (decision === 'always') {
try {
permStore.rememberAlways(name);
} catch (_) {}
decision = 'allow';
}
if (decision !== 'allow') {
const denied = 'permission denied for ' + name;
pushHistory(session, { role: 'tool', name, content: denied, tool_call_id: toolCallId });
emitUpdate(emit, session.id, jobId, { type: 'tool_result', toolCallId, name, result: denied });
continue;
}
}
}
prepared.push({ name, args, toolCallId });
if (name === 'run_terminal_cmd') toolBudget.markShell(budget);
}
let stopEarly = null;
for (const group of toolBatch.groups(prepared)) {
if (cancelled()) return endTurn(emit, session, jobId, tracker, { reason: 'cancelled', turns: turn + 1 });
if (group.sequential) {
for (const item of group.calls) {
stopEarly = await runOneTool(item, turn);
if (stopEarly) break;
}
} else {
const locks = new Map();
const results = await Promise.all(
group.calls.map((item) =>
toolBatch.withPathLock(locks, toolBatch.pathLockKey(item.name, item.args), () => runOneTool(item, turn))
)
);
stopEarly = results.find(Boolean) || null;
}
if (stopEarly) break;
}
if (prepared.length) toolBudget.markToolRound(budget);
if (stopEarly) {
return endTurn(emit, session, jobId, tracker, stopEarly);
}
if (stuckNow) {
return endTurn(emit, session, jobId, tracker, {
reason: 'stuck',
text: lastText || toolBudget.lastToolText(session.history),
turns: turn + 1,
});
}
if (stationarity.shouldNudge(stuck)) {
stationarity.markNudged(stuck);
pushHistory(session, { role: 'user', content: stationarity.nudgeText() });
}
}
return endTurn(emit, session, jobId, tracker, {
reason: 'max_turns',
text: lastText || toolBudget.lastToolText(session.history),
turns: budget.maxTurns,
});
} catch (err) {
if (err && err.message === 'cancelled') {
return endTurn(emit, session, jobId, tracker, { reason: 'cancelled', text: lastText });
}
throw err;
}
}
async function runSubagent(parentCtx, args) {
const prompt = args.prompt || args.description || '';
const label = args.label || 'subagent';
const kind = args.subagent_type === 'general' || args.subagentType === 'general' ? 'general' : 'explore';
let taskId;
try {
taskId = tasks.create(label, { subagentType: kind });
} catch (err) {
return { error: err.message };
}
const rec = tasks.get(taskId);
const hostWorkspace = parentCtx.session.hostWorkspace !== false;
const toolDefs = subagentToolDefs(parentCtx, kind);
await ensureModel(parentCtx.session.model);
const history = [
{
role: 'system',
content: subagentSystem(label, kind, hostWorkspace),
},
{ role: 'user', content: prompt },
];
rec.history = history;
try {
return await runSubagentLoop(parentCtx, rec, history, toolDefs);
} catch (err) {
tasks.fail(taskId, err.message);
throw err;
}
}
function subagentSystem(label, kind, hostWorkspace) {
if (!hostWorkspace) {
return 'You are a focused subagent (' + label + '). You have no host filesystem. Summarize from the prompt only.';
}
if (kind === 'general') {
return (
'You are a focused subagent (' +
label +
'). You may read and write the workspace with the same write restrictions as the parent. Do not spawn further subagents. Summarize when done.'
);
}
return 'You are a focused research subagent (' + label + '). Use read-only tools. Summarize findings.';
}
function subagentToolDefs(parentCtx, kind) {
const hostWorkspace = parentCtx.session.hostWorkspace !== false;
if (!hostWorkspace) return [];
const all = tools.defs({
planMode: planMode.isActive(parentCtx.planTracker),
webFetch: false,
hostWorkspace: true,
});
if (kind !== 'general') {
const allowed = new Set(['read_file', 'grep', 'list_dir', 'memory_search', 'memory_get']);
return all.filter((t) => allowed.has(t.name));
}
const skip = new Set([
'task',
'send_subagent_message',
'wait_tasks',
'kill_task',
'enter_plan_mode',
'exit_plan_mode',
'ask_user_question',
'update_goal',
]);
return all.filter((t) => !skip.has(t.name));
}
async function runSubagentLoop(parentCtx, rec, history, toolDefs) {
let summary = rec.summary || '';
for (let turn = 0; turn < SUBAGENT_TURNS; turn++) {
if (parentCtx.session && live.get(parentCtx.session.id) && live.get(parentCtx.session.id).cancelled) {
throw new Error('cancelled');
}
const ctxSize = loadedCtxSize();
let result;
for (let overflowTry = 0; overflowTry < 4; overflowTry++) {
if (overflowTry > 0 || compaction.shouldCompact(history, toolDefs, ctxSize)) {
const compacted = compaction.compact(history, {
budgetTokens: compaction.historyBudget(ctxSize, toolDefs, overflowTry),
tools: toolDefs,
aggressive: overflowTry > 0,
});
history.length = 0;
for (let i = 0; i < compacted.length; i++) history.push(compacted[i]);
}
try {
result = await engine.complete({
history,
tools: toolDefs,
desktopVision: parentCtx.payload && parentCtx.payload.desktopVision === false ? false : undefined,
});
break;
} catch (err) {
if (!compaction.isOverflowError(err) || overflowTry === 3) throw err;
}
}
if (result.text) history.push({ role: 'assistant', content: result.text });
const calls = result.toolCalls || [];
if (!calls.length) {
summary = result.text || summary;
break;
}
let extra = '';
for (const call of calls) {
try {
const out = await tools.execute(
{
origin: parentCtx.session.origin,
cwd: parentCtx.session.cwd,
session: parentCtx.session,
hostWorkspace: parentCtx.session.hostWorkspace !== false,
planMode: planMode.isActive(parentCtx.planTracker),
planTracker: parentCtx.planTracker,
},
call.name,
typeof call.arguments === 'string' ? safeJson(call.arguments) : call.arguments || {}
);
extra = typeof out === 'string' ? out : JSON.stringify(out);
} catch (err) {
extra = 'error: ' + err.message;
}
history.push({
role: 'tool',
name: call.name,
content: truncate.renderToolResult(extra, toolResultCap()),
tool_call_id: call.id,
});
}
summary = result.text || extra;
tasks.setHistory(rec.id, history);
}
rec.history = history;
tasks.setHistory(rec.id, history);
tasks.finish(rec.id, summary);
return { taskId: rec.id, label: rec.label, summary, subagentType: rec.subagentType || 'explore' };
}
async function continueSubagent(parentCtx, rec, message) {
await ensureModel(parentCtx.session.model);
if (rec.status !== 'running' && tasks.runningCount() >= tasks.MAX_CONCURRENT) {
return { error: 'too many concurrent subagents (max ' + tasks.MAX_CONCURRENT + ')' };
}
const kind = rec.subagentType === 'general' ? 'general' : 'explore';
const hostWorkspace = parentCtx.session.hostWorkspace !== false;
const toolDefs = subagentToolDefs(parentCtx, kind);
let history = Array.isArray(rec.history) && rec.history.length ? rec.history.slice() : null;
if (!history) {
history = [
{ role: 'system', content: subagentSystem(rec.label, kind, hostWorkspace) },
{ role: 'user', content: (rec.summary || '') + '\n\nFollow-up:\n' + message },
];
} else {
history.push({ role: 'user', content: String(message || '') });
}
rec.status = 'running';
rec.history = history;
try {
return await runSubagentLoop(parentCtx, rec, history, toolDefs);
} catch (err) {
tasks.fail(rec.id, err.message);
throw err;
}
}
function safeJson(s) {
try {
return JSON.parse(s);
} catch (_) {
return { raw: s };
}
}
function markLive(sessionId) {
const rec = { cancelled: false };
live.set(sessionId, rec);
if (engine.hold) engine.hold();
return rec;
}
function finishLive(sessionId) {
if (sessionId) live.delete(sessionId);
else live.clear();
if (engine.release) engine.release();
}
function liveCount() {
return live.size;
}
function isLive(sessionId) {
if (sessionId) return live.has(sessionId);
return live.size > 0;
}
function isWaiting() {
return pendingPerms.size + pendingCustom.size + pendingAsks.size + pendingPlans.size > 0;
}
function forget(sessionId) {
if (!sessionId) return;
live.delete(sessionId);
function dropPrefixed(map) {
for (const key of Array.from(map.keys())) {
if (key === sessionId || String(key).indexOf(sessionId) >= 0) {
const rec = map.get(key);
map.delete(key);
if (rec && typeof rec.resolve === 'function') rec.resolve({ error: 'destroyed' });
else if (typeof rec === 'function') rec('deny');
}
}
}
dropPrefixed(pendingCustom);
dropPrefixed(pendingAsks);
dropPrefixed(pendingPlans);
dropPrefixed(pendingPerms);
}
function flushPending(map, payload) {
const keys = Array.from(map.keys());
for (const key of keys) {
const rec = map.get(key);
map.delete(key);
if (rec && typeof rec.resolve === 'function') rec.resolve(payload);
}
}
function cancel(sessionId) {
const rec = live.get(sessionId);
if (rec) rec.cancelled = true;
engine.cancel().catch(() => {});
flushPending(pendingCustom, { error: 'cancelled' });
flushPending(pendingAsks, { error: 'cancelled' });
flushPending(pendingPlans, { decision: 'reject' });
const permKeys = Array.from(pendingPerms.keys());
for (const key of permKeys) {
const fn = pendingPerms.get(key);
pendingPerms.delete(key);
if (typeof fn === 'function') fn('deny');
}
}
module.exports = {
runTurn,
resolvePermission,
resolveCustomResult,
resolveAsk,
resolvePlanDecision,
markLive,
finishLive,
liveCount,
isLive,
isWaiting,
forget,
cancel,
MAX_TURNS,
};