/** * Agentic turn loop: sample (QVAC) → tools → repeat. */ const engine = require('../lib/qvac.js'); const catalog = require('../lib/catalog.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 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 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() { const ctx = loadedCtxSize(); return Math.min(8000, Math.max(1200, Math.floor(ctx * 0.35))); } 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); } } 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; return tools .defs({ planMode: planMode.isActive(tracker), webFetch: payload && payload.webFetch, hostWorkspace, builtinTools: session.builtinTools, }) .concat(customTools.defs(session.id)); } 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 sidecarMessages(session, tracker) { const extra = []; if (tracker && tracker.pendingCompactReminder) { extra.push({ role: 'user', content: compaction.compactReminder() }); const cont = compaction.autoContinue(session.history); if (cont) extra.push(cont); tracker.pendingCompactReminder = false; } if (planMode.isActive(tracker)) { extra.push({ role: 'user', content: '\n' + planMode.reminder(tracker, { hasContent: !!(sessions.readPlan(session.id) || '').trim() }) + '\n', }); } else if (tracker && tracker.pendingExitReminder) { extra.push({ role: 'user', content: '\n' + planMode.exitReminder() + '\n' }); 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: '\n' + note + '\n' }); } 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); let extraSys = payload && payload.system; if (session.goal && goalMod.isActive(session.goal)) { extraSys = [extraSys, goalMod.plannerAddendum(session.goal)].filter(Boolean).join('\n\n'); } const sys = prompts.assemble({ cwd: hostWorkspace ? cwd : session.workspace || cwd, hostWorkspace, extra: extraSys, fsRead: hostWorkspace ? (c, r) => { try { return fsRead(c, r); } catch (_) { return ''; } } : null, }); const sidecars = []; if (hostWorkspace) { 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(); let lastText = ''; let goalNudges = 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()); 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()); 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()); 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 < MAX_TURNS; turn++) { if (cancelled()) return endTurn(emit, session, jobId, tracker, { reason: 'cancelled', turns: turn }); const toolDefs = buildToolDefs(session, payload, tracker); ctx.planMode = planMode.isActive(tracker); const ctxSize = loadedCtxSize(); const beforeLen = session.history.length; 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)) { emitUpdate(emit, session.id, jobId, { type: 'compaction', status: 'start', method: 'llm', used: beforeUsage.used, limit: beforeUsage.limit, pct: beforeUsage.pct, threshold: beforeUsage.threshold, }); session.history = await compaction.compactWithLlm(session.history, { budgetTokens: compaction.historyBudget(ctxSize, toolDefs, 0), tools: toolDefs, complete: (opts) => engine.complete(Object.assign({}, opts, { desktopVision: false })), }); 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)); } else { session.history = compaction.compact(session.history, { budgetTokens: compaction.historyBudget(ctxSize, toolDefs, 0), tools: toolDefs, }); if (session.history.length !== beforeLen) { sessions.replaceHistory(session.id, session.history); const afterUsage = compaction.usage(session.history, toolDefs, ctxSize); emitUpdate(emit, session.id, jobId, { type: 'compaction', status: 'done', method: 'heuristic', 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)); } } emitUpdate(emit, session.id, jobId, { type: 'turn', turn }); let streamBase = compaction.usage(session.history.concat(sidecarMessages(session, tracker)), 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(sidecarMessages(session, tracker)); 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, }, (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, { budgetTokens: compaction.historyBudget(ctxSize, toolDefs, overflowTry + 1), tools: toolDefs, aggressive: true, }); sessions.replaceHistory(session.id, session.history); tracker.pendingCompactReminder = true; emitCompactDone(emit, session, jobId, toolDefs, ctxSize, streamBase, 'overflow'); } } emitLive( emit, session.id, jobId, Object.assign({ type: 'context' }, usageFromStats(result && result.stats, liveUsed(), ctxSize)) ); if (result.text) { lastText = result.text; pushHistory(session, { role: 'assistant', content: result.text }); } const calls = result.toolCalls || []; if (!calls.length) { const goalActive = goalMod.isActive(session.goal); if (goalActive && goalNudges < MAX_GOAL_NUDGES) { goalNudges += 1; pushHistory(session, { role: 'user', content: goalMod.continuation(session.goal) }); continue; } return endTurn(emit, session, jobId, tracker, { type: 'end', reason: 'stop', text: result.text || lastText || '', 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 (sandbox.needsPermission(name, mode) && !planMode.isPlanFilePath(args.path, tracker.planPath)) { 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.remember(name, args, 'allow'); } 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 }); } 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 (stopEarly) { return endTurn(emit, session, jobId, tracker, stopEarly); } if (stuckNow) { return endTurn(emit, session, jobId, tracker, { reason: 'stuck', text: lastText, 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, turns: MAX_TURNS }); } 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, };