// antigravity — use Google Antigravity (agy) models inside the pi coding // agent via the agy stream-json RPC. pi stays the UI: model picker, portable // sessions, and rendering; agy owns native context and runs the selected model with // --dangerously-skip-permissions always enabled (headless agy turns // auto-deny tools that would need a permission prompt otherwise). // // Commands: // /agy show agy conversation and persistent-driver status // /agy-reset drop the current agy conversation (next turn starts fresh) // /agy-models re-discover models from `agy models` and re-register // /agy-agents list configured custom agents // /agy-doctor diagnose binary, models, driver, bridge, and local state // /agy-subagents list subagent activity observed on the agy stream // /agy-usage show Antigravity model quotas (weekly and 5-hour limits) import type { ExtensionAPI, ExtensionContext, ExtensionUIContext, ProviderModelConfig, } from "@earendil-works/pi-coding-agent"; import { getDocsPath, getExamplesPath, getReadmePath } from "@earendil-works/pi-coding-agent"; import { stripTerminalSequences, Text, truncateToWidth } from "@earendil-works/pi-tui"; import { Type } from "typebox"; import { randomUUID } from "node:crypto"; import path from "node:path"; import { piConfigDir, readJson, writeJson } from "./lib/config.ts"; import { getAgyChildrenRegistry, installAgyDeathHooks, killAllAgyTrees, signalAgyTree, signalVerifiedAgyOrphans, } from "./lib/agy-children.ts"; import { checkAgyBinary, MIN_AGY_VERSION, runAgyCommand } from "./lib/agy-diagnostics.ts"; import { parseAgyAgents, readAgyProcessProfile } from "./lib/agy-profile.ts"; import { pruneBridgeMcpCache, removeMcpCacheEntry } from "./lib/mcp-cache.ts"; import { AgyPiBridge, BRIDGE_SERVER_NAME, selectBridgedTools, type PiToolInfo, } from "./lib/bridge.ts"; import { createBridgeLifecycleManager } from "./lib/bridge-lifecycle.ts"; import { ACTIVATE_SKILL_TOOL_NAME, activateSkillDescription, activateSkillParameters, handleActivateSkill, piPrivateSkills, usableSkillCatalog, type SkillLite, } from "./lib/skills.ts"; import { AGY_PI_SCHEDULING_CONTEXT_WINDOW, capabilitiesForModel, FALLBACK_MODELS, mergeAgyModels, modelCacheIsFresh, parseAgyModels, pricingForModel, resolveAgyModelEffort, type AgyModelInfo, } from "./lib/models.ts"; import type { AgyActivity } from "./lib/reducer.ts"; import { AGY_COMPACTION_ENTRY, agyContextTokens, detectAgyCompaction, formatAgyContextTokens, type AgyCompactionMarker, } from "./lib/agy-compaction.ts"; import { AGY_CONVERSATION_STATE_ENTRY, agyConversationExists, restorableAgyConversation, type PersistedAgyConversation, type PersistedAgyReset, } from "./lib/conversation-state.ts"; import { readAgyConversationMetadata } from "./lib/conversation-metadata.ts"; import { AgyReplayStore, type RecordedAgyTool } from "./lib/replay.ts"; import { agyGroupLeadingDescendants, agyGroupSurvivors, agyTaskStopPids, findAgyTask, listAgyTasks, stopAgyTask, type AgyTask, } from "./lib/tasks.ts"; import { formatAgySubagents, trackAgySubagent, type AgySubagentEntry } from "./lib/subagents.ts"; import { findAgyArtifact, listAgyArtifacts } from "./lib/artifacts.ts"; import { fetchAgyUsage } from "./lib/usage.ts"; import { buildAgyRelayedInstructions, WRAPPER_TOOL_DESCRIPTION, WRAPPER_TOOL_NAME, } from "./lib/prompt.ts"; import { wrapperToolActiveAfterModelSwitch } from "./lib/wrapper-activation.ts"; import { openAgyTasksPicker } from "./src/tasks-ui.ts"; import { openArtifact, openAgyArtifactsPicker } from "./src/artifacts-ui.ts"; import { openAgyUsagePicker } from "./src/usage-ui.ts"; import { agyToolLabel, formatAgyCall, summarizeAgyResult } from "./lib/render.ts"; import { mapThinkingToEffort, streamAntigravity } from "./src/provider.ts"; import { AntigravityRuntime, createAntigravityRuntime, runAntigravity, type AntigravityStateSnapshot, } from "./src/runtime.ts"; const MODEL_CACHE_FILE = path.join(piConfigDir("antigravity"), "model-list.json"); const DISCOVERY_TIMEOUT_MS = 15_000; function statusOneLine(value: string, maxLength = 160): string { return stripTerminalSequences(value) .replace(/[\x00-\x1f\x7f]/g, " ") .replace(/\s+/g, " ") .trim() .slice(0, maxLength); } // --- Pi-tool bridge --------------------------------------------------------- const BRIDGE_ENABLED = process.env.PI_ANTIGRAVITY_PI_TOOL_BRIDGE !== "0"; /** Timeout for shutdown-path agy calls — fast enough not to stall closing pi. */ const SHUTDOWN_AGY_TIMEOUT_MS = 5_000; async function execAgy(args: string[], timeoutMs = DISCOVERY_TIMEOUT_MS): Promise { return (await runAgyCommand(args, { timeoutMs })).stdout; } /** * Remove `pi-bridge-*` MCP registrations whose loopback server is no longer * reachable. Registrations live in agy's GLOBAL config while bridge servers * are per-pi-session: crashed sessions leak registrations forever, and every * live session's tools are merged into every agy turn's tools/list. Pruning * dead entries on startup keeps cross-session pollution to live sessions. * agy also caches tool manifests on disk per server * (`~/.gemini/antigravity-cli/mcp//`) and never evicts them, so the * same sweep prunes cache entries with no live server left. */ async function pruneStaleBridgeRegistrations(): Promise { let list: string; try { list = await execAgy(["mcp", "list"]); } catch { return; // agy unavailable — registration below will warn instead } const stale: string[] = []; const live: string[] = []; for (const line of list.split("\n").slice(1)) { const columns = line.trim().split(/\s+/); const name = columns[0] ?? ""; const url = columns[columns.length - 1] ?? ""; if (!name.startsWith(`${BRIDGE_SERVER_NAME}-`) || !/^http:\/\/127\.0\.0\.1:\d+\//.test(url)) { continue; } try { await fetch(url, { method: "POST", headers: { "content-type": "application/json" }, body: JSON.stringify({ jsonrpc: "2.0", id: 1, method: "ping" }), signal: AbortSignal.timeout(1_500), }); // Any HTTP answer (even 403/404) means a server is still listening. live.push(name); } catch { stale.push(name); } } await Promise.all(stale.map((name) => execAgy(["mcp", "remove", name]).catch(() => {}))); // `agy mcp remove` deregisters but leaves the on-disk manifest cache // behind; entries whose registration is already gone never reappear in // `agy mcp list`, so prune the cache against the live set here. await pruneBridgeMcpCache({ liveServers: live }); } interface ModelCache { fetchedAt?: number; source?: "live" | "fallback"; models: AgyModelInfo[]; } /** * Default rates (USD per Mtok) feed pi's native cost calculation. Per-model * overrides belong in pi's own ~/.pi/agent/models.json under * providers.antigravity.modelOverrides — pi applies them over registered models. */ function toProviderModel(model: AgyModelInfo): ProviderModelConfig { const capabilities = capabilitiesForModel(model.id); return { id: model.id, name: model.name, reasoning: true, input: ["text"], cost: pricingForModel(model.id), contextWindow: capabilities.contextWindow, maxTokens: capabilities.maxTokens, }; } /** Collapse effort variants and dedupe (also heals pre-0.2.0 caches). */ function normalizeModels(models: AgyModelInfo[]): AgyModelInfo[] { return mergeAgyModels(models).map((model) => { if (model.supportedEfforts.length > 0) return model; const fallback = FALLBACK_MODELS.find((candidate) => candidate.id === model.id); return fallback?.supportedEfforts.length ? { ...model, supportedEfforts: [...fallback.supportedEfforts], defaultEffort: fallback.defaultEffort, } : model; }); } async function listAgyModels(): Promise { try { return parseAgyModels(await execAgy(["models"])); } catch { return []; } } function getInitialModelCache(): ModelCache { const cached = readJson(MODEL_CACHE_FILE, null); if (cached?.models?.length && modelCacheIsFresh(cached)) { return { ...cached, models: normalizeModels(cached.models) }; } // An expired or missing cache may have been written by an older extension // (e.g. a fallback catalog that predates newer models), so never register // stale disk state at startup: use the catalog baked into this build and // let the startup refresh pull the live list. return { fetchedAt: 0, source: "fallback", models: FALLBACK_MODELS, }; } async function discoverModels(refresh = false): Promise { if (!refresh) { const cached = readJson(MODEL_CACHE_FILE, null); if (cached?.models?.length && modelCacheIsFresh(cached)) { return { ...cached, models: normalizeModels(cached.models) }; } } const live = await listAgyModels(); const cache: ModelCache = live.length ? { fetchedAt: Date.now(), source: "live", models: normalizeModels(live) } : { fetchedAt: Date.now(), source: "fallback", models: FALLBACK_MODELS }; try { writeJson(MODEL_CACHE_FILE, cache); } catch { // Cache is best-effort; discovery still returns. } return cache; } export default function antigravityExtension(pi: ExtensionAPI): void { installAgyDeathHooks(); const runtime = createAntigravityRuntime(); const service = runtime.runSync(AntigravityRuntime); const replay = new AgyReplayStore(); let currentCache = getInitialModelCache(); let observedContextTokens: number | undefined; let persistedConversationKey: string | undefined; // Live subagent roster, folded from stream activities since the last // conversation reset/restore. In-memory only: agy's own conversation // database remains the durable record. const subagentRoster = new Map(); let selectedModelKey: string | undefined; pi.registerEntryRenderer(AGY_COMPACTION_ENTRY, (entry, _options, theme) => { const marker = entry.data as AgyCompactionMarker; const text = theme.fg("dim", "── ") + theme.fg("accent", "agy compacted context") + theme.fg( "dim", ` · ~${formatAgyContextTokens(marker.beforeTokens)} → ~${formatAgyContextTokens(marker.afterTokens)} ──`, ); return { render: (width: number) => [truncateToWidth(text, width, "")], invalidate: () => {}, }; }); const conversationStateKey = (state: { conversationId: string; modelId: string; turns: number; }): string => `${state.conversationId}:${state.modelId}:${state.turns}`; /** * Every session-bound getter on a stale ctx throws once pi invalidates the * extension runner (dispose/newSession/fork/switchSession/reload), and * events dispatched around the swap — e.g. agent_settled landing after a * /new or a shutdown — still reach our handlers. Probe first; a dead session * has nothing to persist into and no widget to update. */ function probeCtx(read: () => T): { stale: true } | { stale: false; value: T } { try { return { stale: false, value: read() }; } catch { return { stale: true }; } } async function persistConversationState(ctx: ExtensionContext, force = false): Promise { // Read every session-bound field before the first await — the stale window // widens once this handler yields. const probe = probeCtx(() => ({ provider: ctx.model?.provider, sessionId: ctx.sessionManager.getSessionId(), cwd: ctx.cwd, })); if (probe.stale || probe.value.provider !== "antigravity") return; const { sessionId, cwd } = probe.value; const snapshot = await runAntigravity(runtime, service.snapshot); if (!snapshot.conversationId || !snapshot.model || snapshot.cwd !== cwd) return; const state: PersistedAgyConversation = { version: 1, kind: "conversation", sessionId, conversationId: snapshot.conversationId, cwd, modelId: snapshot.model, turns: snapshot.turns, usage: snapshot.conversationUsage, contextTokens: observedContextTokens, }; const key = conversationStateKey(state); if (!force && key === persistedConversationKey) return; try { pi.appendEntry(AGY_CONVERSATION_STATE_ENTRY, state); } catch { return; // Session went stale mid-persist. } persistedConversationKey = key; } function appendConversationReset(ctx: ExtensionContext): void { const probe = probeCtx(() => ({ sessionId: ctx.sessionManager.getSessionId(), cwd: ctx.cwd, })); if (probe.stale) return; const reset: PersistedAgyReset = { version: 1, kind: "reset", sessionId: probe.value.sessionId, cwd: probe.value.cwd, }; try { pi.appendEntry(AGY_CONVERSATION_STATE_ENTRY, reset); } catch { return; } observedContextTokens = undefined; persistedConversationKey = undefined; } // --- Pi-tool bridge setup ------------------------------------------------- const bridge = new AgyPiBridge(`${BRIDGE_SERVER_NAME}-${process.pid}`); const bridgeToken = randomUUID(); bridge.requireToken(bridgeToken); // Session-unique tool prefix: agy's MCP config is global, so concurrent pi // sessions' bridges all appear in every agy turn's tools/list. Namespacing // each session's tools makes tool→server routing unambiguous — a call can // only reach the session that advertised it. const bridgeToolPrefix = `pi__p${process.pid}__`; bridge.setToolPrefix(bridgeToolPrefix); bridge.setOnCall((call) => service.pushBridgeCall(call)); const bridgeManager = createBridgeLifecycleManager({ bridge, bridgeToken, enabled: BRIDGE_ENABLED, pruneStaleRegistrations: pruneStaleBridgeRegistrations, addMcpServer: async (serverName, url, token) => { await execAgy([ "mcp", "add", "--type", "http", "--header", `x-pi-bridge-token: ${token}`, serverName, url, ]); }, removeMcpServer: async (serverName) => { await execAgy(["mcp", "remove", serverName], SHUTDOWN_AGY_TIMEOUT_MS); }, evictMcpCache: removeMcpCacheEntry, // tasksUi is populated by the session_start handler; read it lazily so a // warning raised before UI attach is simply dropped rather than throwing. notifyWarning: (message) => tasksUi?.notify(message, "warning"), }); /** * Register the pi-tool bridge with agy. Idempotent: no-op when disabled or * already registered. Registration must precede the first agy spawn of the * session — agy eagerly connects to MCP servers at startup (verified * 2026-08-21) — so this runs as soon as an Antigravity model is selected. */ async function ensureBridgeRegistered(ui?: ExtensionUIContext): Promise { await bridgeManager.ensureRegistered((warning) => ui?.notify(warning, "warning")); } /** * Deregister the pi-tool bridge and evict its manifest cache. No-op when * the bridge is not registered or running. Runs when the session leaves * Antigravity models and on shutdown. */ async function teardownBridge(): Promise { await bridgeManager.teardown(); } // --- Skill passing (Phase 2) ---------------------------------------------- // pi's loaded skills, refreshed per turn via before_agent_start so /reload // is respected. Model-invocation-disabled skills are excluded. let loadedSkills: SkillLite[] = []; /** * The instruction block relayed to agy on fresh conversations: pi's * documentation section only — boilerplate, tool inventory, skills and * workspace rules all reach agy through other channels or not at all. */ const getSystemPromptRelay = () => buildAgyRelayedInstructions({ readme: getReadmePath(), docs: getDocsPath(), examples: getExamplesPath(), }); function captureSkills(skills: unknown): void { if (!Array.isArray(skills)) return; loadedSkills = skills .map((skill) => skill as Partial & { disableModelInvocation?: boolean }) .filter( (skill) => skill.disableModelInvocation !== true && typeof skill.filePath === "string", ) .map((skill) => ({ name: String(skill.name), description: String(skill.description ?? ""), filePath: String(skill.filePath), baseDir: String(skill.baseDir ?? path.dirname(String(skill.filePath))), })); } const bridgedSkills = () => usableSkillCatalog(piPrivateSkills(loadedSkills, tasksSessionCwd)); /** * Bridge mode keeps the catalog in activate_skill's schema (refreshed on * every agy spawn), so nothing is appended to the prompt. When the bridge * is off OR failed to register, pi-private skills are simply unavailable — * warn the user once instead of stuffing the catalog into the prompt. * The undefined return is intentional; the runtime's bootstrapSuffix * plumbing stays so a future direct-mode catalog needs no new wiring. */ const getBootstrapSuffix = () => { bridgeManager.warnSkillsUnavailable(bridgedSkills()); return undefined; }; /** Publish one `pi__p__activate_skill` tool for global pi skills. */ function refreshSkillTools(): void { if (!BRIDGE_ENABLED) return; const skills = bridgedSkills(); if (skills.length === 0) { bridge.setDynamicTools([]); return; } bridge.setDynamicTools([ { name: ACTIVATE_SKILL_TOOL_NAME, description: activateSkillDescription(skills), parameters: activateSkillParameters(skills), handler: (args) => handleActivateSkill(skills, args), }, ]); } bridge.setToolSource(() => { if (!BRIDGE_ENABLED) return []; let activeNames: string[] = []; let allTools: PiToolInfo[] = []; try { activeNames = pi.getActiveTools(); allTools = pi.getAllTools() as PiToolInfo[]; } catch { return []; // API unavailable (print/RPC edge) — expose nothing. } return selectBridgedTools(allTools, activeNames); }); // Per-skill tools are gone; one activate_skill tool is published on every // skills capture via refreshSkillTools(). // --- Status-bar hint for live agy background tasks ----------------------- const AGY_TASKS_WIDGET_KEY = "agy-tasks"; const AGY_ARTIFACTS_WIDGET_KEY = "agy-artifacts"; let tasksUi: ExtensionUIContext | undefined; let tasksSessionCwd: string | undefined; let widgetLiveCount = -1; let widgetArtifactCount = -1; let widgetScanInFlight = false; let widgetScanQueued = false; let agyAgentActive = false; let widgetPollTimer: ReturnType | undefined; const WIDGET_POLL_MS = 2_000; /** Active-tool snapshot for the provider's native re-execution fallback. */ const isActiveTool = (name: string): boolean => { try { return pi.getActiveTools().includes(name); } catch { return false; // API unavailable (print/RPC edge) — stay on the wrapper. } }; function setAgyTasksWidget(live: number): void { if (!tasksUi || live === widgetLiveCount) return; widgetLiveCount = live; try { if (live === 0) { tasksUi.setWidget(AGY_TASKS_WIDGET_KEY, undefined); return; } tasksUi.setWidget(AGY_TASKS_WIDGET_KEY, (_tui, theme) => { const line = theme.fg("warning", "■ ") + theme.fg("text", `${live} agy background task${live === 1 ? "" : "s"}`) + theme.fg("dim", " • ") + theme.fg("accent", "/agy-tasks") + theme.fg("dim", " to view"); return { render: (width: number) => [truncateToWidth(line, width, "")], invalidate: () => {}, }; }); } catch { // UI may be unavailable (print/RPC modes or teardown). } } function setAgyArtifactsWidget(count: number): void { if (!tasksUi || count === widgetArtifactCount) return; widgetArtifactCount = count; try { if (count === 0) { tasksUi.setWidget(AGY_ARTIFACTS_WIDGET_KEY, undefined); return; } tasksUi.setWidget(AGY_ARTIFACTS_WIDGET_KEY, (_tui, theme) => { const line = theme.fg("success", "◆ ") + theme.fg("text", `${count} agy artifact${count === 1 ? "" : "s"}`) + theme.fg("dim", " • ") + theme.fg("accent", "/agy-artifacts") + theme.fg("dim", " to view"); return { render: (width: number) => [truncateToWidth(line, width, "")], invalidate: () => {}, }; }); } catch { // UI may be unavailable (print/RPC modes or teardown). } } function reconcileWidgetPolling(): void { const shouldPoll = Boolean(tasksUi && (agyAgentActive || widgetLiveCount > 0)); if (shouldPoll && !widgetPollTimer) { widgetPollTimer = setInterval(updateAgyTasksWidget, WIDGET_POLL_MS); widgetPollTimer.unref?.(); } else if (!shouldPoll && widgetPollTimer) { clearInterval(widgetPollTimer); widgetPollTimer = undefined; } } /** Rescan task state independently from agy's provider stream. */ function updateAgyTasksWidget(): void { if (!tasksUi) return; if (widgetScanInFlight) { widgetScanQueued = true; return; } widgetScanInFlight = true; void (async () => { try { const snapshot = await runAntigravity(runtime, service.snapshot); if (!snapshot.conversationId) { setAgyTasksWidget(0); setAgyArtifactsWidget(0); return; } const [tasks, artifacts] = await Promise.all([ listAgyTasks(snapshot.conversationId, { sessionCwd: tasksSessionCwd, agyPids: [...getAgyChildrenRegistry().live], }), listAgyArtifacts(snapshot.conversationId), ]); setAgyTasksWidget( tasks.filter( (task) => task.pids.length > 0 || task.orphans.length > 0 || task.ambiguous.length > 0, ).length, ); setAgyArtifactsWidget(artifacts.length); } catch { // Runtime closed or scan failed; leave the widget as-is. } finally { widgetScanInFlight = false; if (widgetScanQueued) { widgetScanQueued = false; queueMicrotask(updateAgyTasksWidget); } else { reconcileWidgetPolling(); } } })(); } function handleAgyActivity(activity: AgyActivity): void { if (activity.type === "conversation_fallback") { observedContextTokens = undefined; persistedConversationKey = undefined; subagentRoster.clear(); return; } trackAgySubagent(subagentRoster, activity); if (activity.type === "usage") { const nextContextTokens = agyContextTokens(activity.usage); const compaction = detectAgyCompaction(observedContextTokens, nextContextTokens); if (compaction) { const marker: AgyCompactionMarker = { version: 1, ...compaction, detectedAt: new Date().toISOString(), }; try { pi.appendEntry(AGY_COMPACTION_ENTRY, marker); } catch { // Session replaced mid-activity; the marker is advisory only. } } observedContextTokens = nextContextTokens; return; } if ( (activity.type === "tool_start" || activity.type === "tool_done" || activity.type === "tool_error") && (activity.name === "run_command" || activity.name === "schedule") ) { // ACTIVE arrives before a sleeping command completes. Start the // independent filesystem scan now instead of waiting for onSettled, // then retry once in case the task log and holder fd are still racing. updateAgyTasksWidget(); if (activity.type === "tool_start") { const retry = setTimeout(updateAgyTasksWidget, 500); retry.unref?.(); } } } pi.registerTool({ name: WRAPPER_TOOL_NAME, label: "antigravity", description: WRAPPER_TOOL_DESCRIPTION, parameters: Type.Object({ tool: Type.String(), input: Type.Unknown(), }), async execute(toolCallId, params) { const recorded = replay.take(toolCallId); if (!recorded) { throw new Error(`No recorded antigravity result for "${params.tool}".`); } if (recorded.error) { throw new Error(recorded.error); } const body = recorded.output ?? ""; return { content: [{ type: "text", text: body ? body.slice(0, 16_000) : "(no output)" }], details: recorded, }; }, renderCall(args, theme) { return new Text(formatAgyCall(args.tool, args.input, theme), 0, 0); }, renderResult(result, { expanded }, theme, context) { const body = result.content[0]?.type === "text" ? result.content[0].text : ""; const details = result.details as RecordedAgyTool | undefined; const tool = details?.agyTool ?? "tool"; if (context.isError) { const message = body && body !== "(no output)" ? body.split("\n")[0] : "failed"; return new Text(theme.fg("error", `✗ ${agyToolLabel(tool)}: ${message}`), 0, 0); } const secs = typeof details?.durationSeconds === "number" ? theme.fg("muted", ` (${details.durationSeconds.toFixed(2)}s)`) : ""; const { counts } = summarizeAgyResult(tool, details?.output); const parts = [theme.fg("success", "✓ "), counts ? theme.fg("muted", counts) : "", secs]; let text = parts.join(""); if (body && body !== "(no output)") { const lines = body.split("\n"); const shown = expanded ? lines : lines.slice(0, 3); text += `\n${shown.map((line) => theme.fg("toolOutput", line)).join("\n")}`; if (!expanded && lines.length > 3) { text += theme.fg("muted", `\n… +${lines.length - 3} lines (ctrl+o to expand)`); } } return new Text(text, 0, 0); }, }); const registerAntigravityProvider = (models: AgyModelInfo[]) => { pi.registerProvider("antigravity", { name: "Google Antigravity (agy)", baseUrl: "agy://local-stream-json", apiKey: "agy-local-session", api: "antigravity-stream-json", models: models.map(toProviderModel), streamSimple: streamAntigravity( runtime, service, replay, bridge, updateAgyTasksWidget, getBootstrapSuffix, isActiveTool, handleAgyActivity, createAntigravityRuntime, (modelId) => currentCache.models.find((candidate) => candidate.id === modelId), readAgyProcessProfile, bridgeManager.processRevision, getSystemPromptRelay, ), }); }; registerAntigravityProvider(currentCache.models); async function refreshStaleModelsWhenSelected(): Promise { if (modelCacheIsFresh(currentCache)) return; const fresh = await discoverModels(true); currentCache = fresh; registerAntigravityProvider(fresh.models); } // Startup registration must stay synchronous, and the picker needs a // correct list even when the default model is not from antigravity (the // session_start/model_select hooks only refresh once agy is selected): a // no-op for fresh caches, otherwise a background heal of the registration. void refreshStaleModelsWhenSelected().catch((error) => { console.error("pi-antigravity: background model refresh failed", error); }); pi.on("before_agent_start", (event) => { captureSkills(event.systemPromptOptions?.skills); refreshSkillTools(); }); // The agy stream reports tool starts immediately, but its DONE event can be // delayed by a sleeping command. Poll the independent filesystem task // source while the agent or any discovered task is live, then stop when // both are idle. This keeps the widget current without a permanent timer. pi.on("agent_start", (_event, ctx) => { const probe = probeCtx(() => ctx.model?.provider); if (probe.stale || probe.value !== "antigravity") return; agyAgentActive = true; updateAgyTasksWidget(); reconcileWidgetPolling(); }); pi.on("agent_settled", async (_event, ctx: ExtensionContext) => { const probe = probeCtx(() => ctx.model?.provider); if (!probe.stale && (probe.value === "antigravity" || agyAgentActive)) { agyAgentActive = false; updateAgyTasksWidget(); } await persistConversationState(ctx); }); // The display-only `antigravity` wrapper tool only matters while an agy // model is active (the provider synthesizes its toolCalls from recorded // agy activity). Keep it out of every other model's tool payload: sync // active-tool state whenever the selected model changes, including the // session's initial restore. const syncWrapperToolActivation = (provider: string | undefined) => { try { const next = wrapperToolActiveAfterModelSwitch( pi.getActiveTools(), WRAPPER_TOOL_NAME, provider, "antigravity", ); if (next) pi.setActiveTools([...next]); } catch { // Tool APIs unavailable (print/RPC edge) — leave the set untouched. } }; pi.on("session_start", async (event, ctx: ExtensionContext) => { syncWrapperToolActivation(ctx.model?.provider); selectedModelKey = ctx.model ? `${ctx.model.provider}:${ctx.model.id}` : undefined; const restored = event.reason !== "fork" && ctx.model?.provider === "antigravity" ? restorableAgyConversation( ctx.sessionManager.getBranch(), ctx.sessionManager.getSessionId(), ctx.cwd, ) : undefined; const compatibleRestore = restored && restored.modelId === ctx.model?.id && (await agyConversationExists(restored.conversationId)) ? restored : undefined; await runAntigravity( runtime, service.setSession(ctx.cwd, undefined, !compatibleRestore && event.reason !== "new"), ); if (compatibleRestore) { await runAntigravity( runtime, service.restoreConversation({ conversationId: compatibleRestore.conversationId, modelId: compatibleRestore.modelId, cwd: compatibleRestore.cwd, turns: compatibleRestore.turns, usage: compatibleRestore.usage, }), ); observedContextTokens = compatibleRestore.contextTokens; persistedConversationKey = conversationStateKey(compatibleRestore); } else { observedContextTokens = undefined; persistedConversationKey = undefined; } if (ctx.hasUI) tasksUi = ctx.ui; tasksSessionCwd = ctx.cwd; updateAgyTasksWidget(); // The pi-tool bridge is only useful while an Antigravity model is // selected: register lazily here (session resumed on an agy model) and on // model_select; non-agy sessions never touch agy at all. if (ctx.model?.provider === "antigravity") { await ensureBridgeRegistered(ctx.ui); await refreshStaleModelsWhenSelected(); } }); pi.on("model_select", async (event, ctx) => { syncWrapperToolActivation(event.model?.provider); const nextModelKey = event.model ? `${event.model.provider}:${event.model.id}` : undefined; if (selectedModelKey?.startsWith("antigravity:") && selectedModelKey !== nextModelKey) { // Another provider/model can add context that the mutable agy // conversation never saw. Force a branch bootstrap when agy is selected // again instead of silently resuming stale native history. const cwd = probeCtx(() => ctx.cwd); if (!cwd.stale) { await runAntigravity(runtime, service.setSession(cwd.value, undefined, true)); observedContextTokens = undefined; persistedConversationKey = undefined; } } selectedModelKey = nextModelKey; // The bridge exists only while an Antigravity model is selected. if (event.model?.provider === "antigravity") { const ui = probeCtx(() => ctx?.ui); await ensureBridgeRegistered(ui.stale ? undefined : ui.value); await refreshStaleModelsWhenSelected(); } else { await teardownBridge(); } }); pi.on("session_tree", async (_event, ctx: ExtensionContext) => { // An agy conversation cannot be rewound to match a different pi branch. // Restart it and bootstrap the selected branch on the next provider call. const cwd = probeCtx(() => ctx.cwd); if (cwd.stale) return; await runAntigravity(runtime, service.setSession(cwd.value, undefined, true)); appendConversationReset(ctx); subagentRoster.clear(); setAgyTasksWidget(0); setAgyArtifactsWidget(0); }); pi.on("session_compact", async (_event, ctx: ExtensionContext) => { // Pi manual/overflow compaction changes the session branch but not agy's // native conversation. Re-anchor the same owner after the compaction entry. await persistConversationState(ctx, true); }); pi.on("session_shutdown", async () => { agyAgentActive = false; if (widgetPollTimer) clearInterval(widgetPollTimer); widgetPollTimer = undefined; widgetScanQueued = false; // Group ids signalled at the top of shutdown; survivors get SIGKILL after // service.close has given SIGTERM its grace. Every target in this set was // verified by live ancestry or a current lsof hold moments earlier. const sweepTargets = new Set(); // Signal every process-group-leading child of our tracked agy processes. // agy >= 1.2.0 runs tasks in their own process groups, so killAllAgyTrees // (which only reaches agy's own group) leaves them running; ancestry is // proof of ownership here — no per-task attribution needed — while agy is // still alive to be their parent. try { const leaders = await agyGroupLeadingDescendants([...getAgyChildrenRegistry().live]); for (const pid of leaders) { sweepTargets.add(pid); signalAgyTree({ pid }, "SIGTERM"); } } catch { // Process scan failed; service.close and killAllAgyTrees still run. } // Orphans recorded at earlier recycles are proven descendants of agy // processes this pi spawned — ours regardless of which conversation the // service snapshot currently names (or whether a /agy reset cleared it). // signalVerifiedAgyOrphans re-checks each record's identity first, so a // stale record can never signal a reused pid or a stranger's group. signalVerifiedAgyOrphans("SIGTERM"); // Stop any live agy background tasks so closing pi leaves nothing // running silently. A task stop only ever signals proven log holders; // recorded orphans were covered above, and advisory matches are never // signalled. try { const snapshot = await runAntigravity(runtime, service.snapshot); if (snapshot.conversationId) { const tasks = await listAgyTasks(snapshot.conversationId, { sessionCwd: tasksSessionCwd, }); const live = tasks.filter((task) => task.pids.length > 0); const stopped = await Promise.all(live.map((task) => stopAgyTask(task))); for (const { pgids } of stopped) { for (const pgid of pgids) sweepTargets.add(pgid); } } } catch { // Runtime closed or scan failed; nothing to stop. } // Close the runtime: aborts any in-flight agy child process, then tear // down the Effect runtime. try { await runAntigravity(runtime, service.close); } catch { // Already closed. } if (bridgeManager.isRegistered() || bridgeManager.isRunning()) { await teardownBridge(); } try { tasksUi?.setWidget(AGY_TASKS_WIDGET_KEY, undefined); tasksUi?.setWidget(AGY_ARTIFACTS_WIDGET_KEY, undefined); } catch { // UI may already be gone. } tasksUi = undefined; widgetLiveCount = -1; widgetArtifactCount = -1; widgetScanInFlight = false; try { await runtime.dispose(); } catch { // Disposed gracefully } // Escalate: any signalled group that still holds members ignored its // SIGTERM — force it now. -pgid addresses the group, so a dead leader's // pid can never be confused with a reused process. Recorded orphans go // through signalVerifiedAgyOrphans again — identity is re-checked, never // inherited from the earlier sweep. try { for (const pgid of await agyGroupSurvivors(sweepTargets)) { try { process.kill(-pgid, "SIGKILL"); } catch { // Group already gone. } } } catch { // Escalation scan failed; killAllAgyTrees still runs. } signalVerifiedAgyOrphans("SIGKILL"); // Sweep any remaining tracked agy process trees (including earlier turns // that finished logically while grandchildren held stdio open). killAllAgyTrees(); }); pi.registerCommand("agy", { description: "Show agy conversation and persistent-driver status", handler: async (args, ctx) => { if (args.trim()) { ctx.ui.notify( 'antigravity: use "/agy-reset", "/agy-models", "/agy-agents", "/agy-doctor", "/agy-subagents", "/agy-tasks", "/agy-artifacts", or "/agy-usage".', "error", ); return; } const snapshot = await runAntigravity(runtime, service.snapshot); const id = snapshot.conversationId; const metadata = id ? await readAgyConversationMetadata(id) : undefined; const metadataTitle = metadata?.metadata?.title ? statusOneLine(metadata.metadata.title) : undefined; const title = metadataTitle || (id ? id.slice(0, 12) : "none — next turn starts fresh"); const updated = metadata?.metadata?.updatedAt ? new Date(metadata.metadata.updatedAt).toLocaleString() : undefined; const details = [ `model: ${snapshot.model ?? "unselected"}`, `turns: ${snapshot.turns}`, `driver: ${snapshot.executor.mode}/${snapshot.executor.state}${snapshot.executor.pid ? ` pid=${snapshot.executor.pid}` : ""}`, observedContextTokens === undefined ? undefined : `native context: ~${formatAgyContextTokens(observedContextTokens)}/${formatAgyContextTokens(AGY_PI_SCHEDULING_CONTEXT_WINDOW)}`, metadata?.metadata?.numSteps === undefined ? undefined : `native steps: ${metadata.metadata.numSteps}`, updated ? `updated: ${updated}` : undefined, ].filter((part): part is string => part !== undefined); let profile = "agent: none · mode: default"; try { const configured = readAgyProcessProfile(); profile = `agent: ${configured.agent ?? "none"} · mode: ${configured.mode ?? "default"}`; } catch (error) { profile = `profile error: ${error instanceof Error ? error.message : String(error)}`; } ctx.ui.notify( `antigravity: ${title}\nconversation: ${id ?? "none"}\n${details.join(" · ")}\n${profile}`, "info", ); }, }); pi.registerCommand("agy-reset", { description: "Drop the agy conversation and driver; next turn starts fresh", handler: async (args, ctx) => { if (args.trim()) { ctx.ui.notify('agy-reset: usage "/agy-reset".', "error"); return; } await runAntigravity(runtime, service.reset); subagentRoster.clear(); appendConversationReset(ctx); ctx.ui.notify("antigravity: conversation reset; next turn starts fresh.", "info"); }, }); pi.registerCommand("agy-models", { description: "Re-discover models from `agy models` and re-register the provider", handler: async (args, ctx) => { if (args.trim()) { ctx.ui.notify('agy-models: usage "/agy-models".', "error"); return; } const refreshed = await discoverModels(true); currentCache = refreshed; registerAntigravityProvider(refreshed.models); ctx.ui.notify( `antigravity: ${refreshed.models.length} models registered (${refreshed.source}).`, "info", ); }, }); pi.registerCommand("agy-agents", { description: "List configured custom agy agents", handler: async (args, ctx) => { if (args.trim()) { ctx.ui.notify('agy-agents: usage "/agy-agents".', "error"); return; } try { const agents = parseAgyAgents(await execAgy(["agents"])); ctx.ui.notify( agents.length > 0 ? `antigravity custom agents:\n${agents.map((agent) => `• ${agent}`).join("\n")}` : "antigravity: no custom agents configured.", "info", ); } catch (error) { ctx.ui.notify( `antigravity: failed to list agents (${error instanceof Error ? error.message : String(error)}).`, "error", ); } }, }); pi.registerCommand("agy-doctor", { description: "Diagnose binary, models, driver, bridge, and conversation state", handler: async (args, ctx) => { if (args.trim()) { ctx.ui.notify('agy-doctor: usage "/agy-doctor".', "error"); return; } const lines = ["antigravity doctor"]; const binary = await checkAgyBinary({ refresh: true }); if (binary.ok) { lines.push( `binary: ${binary.binary} (${binary.version}, ${binary.source})`, `binary selection: ${binary.selectionReason ?? "compatible candidate"}`, `minimum: ${MIN_AGY_VERSION}`, ); } else { lines.push( `binary: ERROR [${binary.category}] ${binary.message}`, `minimum: ${MIN_AGY_VERSION}`, ); } for (const candidate of binary.candidates ?? []) { const selected = binary.ok && candidate.binary === binary.binary ? "selected" : "candidate"; lines.push( `binary ${selected}: ${candidate.source} ${candidate.binary} · ${ candidate.ok ? candidate.development ? (candidate.version ?? "development") : (candidate.version ?? "unknown") : `ERROR [${candidate.category ?? "unknown"}]` }`, ); } try { const discovered = parseAgyModels((await runAgyCommand(["models"])).stdout); if (discovered.length === 0) throw new Error("no valid model rows returned"); currentCache = { fetchedAt: Date.now(), source: "live", models: discovered }; registerAntigravityProvider(discovered); lines.push(`models: ${discovered.length} (live)`); } catch (error) { lines.push( `models: ERROR ${error instanceof Error ? error.message : String(error)}; ${currentCache.models.length} cached (${currentCache.source})`, ); } let snapshot: AntigravityStateSnapshot | undefined; try { snapshot = await runAntigravity(runtime, service.snapshot); const executor = snapshot.executor; const config = executor.config; lines.push( `driver: ${executor.mode} · ${executor.state}${executor.pid ? ` · pid ${executor.pid}` : ""}`, `driver binary: ${config?.binary ?? "none"}${config?.binaryVersion ? ` · ${config.binaryVersion}` : ""}`, `driver config: model=${config?.model ?? "none"} effort=${config?.effort ?? "none"} agent=${config?.agent ?? "none"} mode=${config?.mode ?? "default"}`, ); if (executor.stats) { const stats = executor.stats; const reasons = Object.entries(stats.recycleReasons) .map(([reason, count]) => `${reason}=${count}`) .join(", "); lines.push( executor.mode === "persistent" ? `driver stats: spawns=${stats.spawnCount} respawns=${Math.max(0, stats.spawnCount - 1)} turns=${stats.submittedTurns} reused=${stats.reusedTurns} recycles=${stats.recycleCount} current=${stats.currentProcessTurns}` : `driver stats: one-shot launches=${stats.spawnCount} turns=${stats.submittedTurns}`, `driver recycle reasons: ${reasons || "none"}`, ); } } catch (error) { lines.push(`driver: ERROR ${error instanceof Error ? error.message : String(error)}`); } const selected = currentCache.models.find((model) => model.id === ctx.model?.id); let profile = "agent=none mode=default"; try { const configured = readAgyProcessProfile(); profile = `agent=${configured.agent ?? "none"} mode=${configured.mode ?? "default"}`; lines.push( `selection: ${ctx.model?.id ?? "none"} · effort ${ resolveAgyModelEffort( selected, ctx.thinkingLevel === "off" ? undefined : mapThinkingToEffort( ctx.thinkingLevel as Exclude, ), ) ?? "none" }`, ); } catch (error) { lines.push(`profile: ERROR ${error instanceof Error ? error.message : String(error)}`); } lines.push(`profile: ${profile}`); lines.push( `bridge: enabled=${BRIDGE_ENABLED} running=${bridgeManager.isRunning()} registered=${bridgeManager.isRegistered()} revision=${bridgeManager.processRevision()}`, ); const conversationId = snapshot?.conversationId; if (conversationId) { const [exists, metadata] = await Promise.all([ agyConversationExists(conversationId), readAgyConversationMetadata(conversationId), ]); lines.push( `conversation: ${conversationId} · db ${exists ? "readable" : "missing/unreadable"}`, `metadata: ${metadata.status}`, ); } else { lines.push("conversation: none", "metadata: not applicable"); } ctx.ui.notify(lines.join("\n"), binary.ok ? "info" : "error"); }, }); pi.registerCommand("agy-tasks", { description: "List agy background tasks; `stop |stop all` to terminate", handler: async (args, ctx) => { const arg = args.trim().toLowerCase(); const snapshot = await runAntigravity(runtime, service.snapshot); const conversationId = snapshot.conversationId; // No arguments: interactive dashboard overlay (x stops, r rescans). if (!arg) { if (!conversationId) { ctx.ui.notify("agy-tasks: no agy conversation in this session yet.", "error"); return; } const rescan = () => listAgyTasks(conversationId, { sessionCwd: ctx.cwd, agyPids: [...getAgyChildrenRegistry().live], }); await openAgyTasksPicker(ctx, rescan); updateAgyTasksWidget(); return; } const stopMatch = arg.match(/^stop\s+(.+)$/); if (!stopMatch) { ctx.ui.notify('agy-tasks: usage "/agy-tasks" or "/agy-tasks stop |all".', "error"); return; } const target = stopMatch[1].trim(); // `stop all` reaps recorded groups whose agy parent has exited — scoped // to the identity-verified registry, not to the task list, so missing // or deleted task logs never block it. With no live conversation (a // /agy reset cleared the snapshot, not the registry) the sweep runs // unscoped: every recorded group is this pi's own work either way. // protectAttached keeps still-attached groups safe in both cases — // they may be a foreground command under a running driver. Per-task // stops can never do this: nothing binds a recorded process to a task. const orphanGroups = target === "all" ? signalVerifiedAgyOrphans("SIGTERM", { conversationId, protectAttached: true, }) : 0; if (!conversationId) { ctx.ui.notify( orphanGroups > 0 ? `agy-tasks: reaped ${orphanGroups} recorded orphan group(s).` : "agy-tasks: no agy conversation in this session yet.", orphanGroups > 0 ? "info" : "error", ); if (orphanGroups > 0) updateAgyTasksWidget(); return; } // The task listing is best-effort — a log disappearing mid-scan or a // failed `ps` must not block the verified orphan sweep above. let tasks: AgyTask[] = []; try { tasks = await listAgyTasks(conversationId, { sessionCwd: ctx.cwd, agyPids: [...getAgyChildrenRegistry().live], }); } catch (error) { ctx.ui.notify( `agy-tasks: task scan failed (${error instanceof Error ? error.message : String(error)}).`, "warning", ); } const selected = target === "all" ? tasks.filter( (task) => task.pids.length > 0 || task.orphans.length > 0 || task.ambiguous.length > 0, ) : [findAgyTask(tasks, target)].filter( (task): task is NonNullable => task !== undefined, ); if (selected.length === 0) { ctx.ui.notify( orphanGroups > 0 ? `agy-tasks: no running task "${target}" in this conversation; reaped ${orphanGroups} recorded orphan group(s).` : `agy-tasks: no running task "${target}" in this conversation.`, orphanGroups > 0 ? "info" : "error", ); if (orphanGroups > 0) updateAgyTasksWidget(); return; } const stoppable = selected.filter((task) => agyTaskStopPids(task).length > 0); if (stoppable.length === 0 && orphanGroups === 0) { ctx.ui.notify( `agy-tasks: ${selected.map((task) => task.id).join(", ")} has no provably-owned process (unclear); nothing signalled. Check the process manually before killing.`, "warning", ); return; } const results = await Promise.all(stoppable.map((task) => stopAgyTask(task))); const stopped = stoppable.map((task) => task.id).join(", "); const skipped = selected.length - stoppable.length; const signalled = results.reduce((sum, result) => sum + result.signaled, 0) + orphanGroups; ctx.ui.notify( `agy-tasks: sent SIGTERM to ${signalled} process(es)${stopped ? ` via ${stopped}` : ""}${orphanGroups > 0 ? `, incl. ${orphanGroups} recorded orphan group(s)` : ""}${skipped > 0 ? `; ${skipped} skipped as unclear` : ""}.`, "info", ); updateAgyTasksWidget(); }, }); pi.registerCommand("agy-artifacts", { description: "List the agy conversation's artifacts (agent-created files, uploads)", handler: async (args, ctx) => { const snapshot = await runAntigravity(runtime, service.snapshot); const conversationId = snapshot.conversationId; if (!conversationId) { ctx.ui.notify("agy-artifacts: no agy conversation in this session yet.", "error"); return; } const rescan = () => listAgyArtifacts(conversationId); const arg = args.trim(); // `open `: non-interactive open by exact name or unique prefix. const openMatch = arg.match(/^open\s+(.+)$/); if (openMatch) { const artifacts = await rescan(); const artifact = findAgyArtifact(artifacts, openMatch[1]); if (!artifact) { ctx.ui.notify(`agy-artifacts: no artifact matching "${openMatch[1]}".`, "error"); return; } try { await openArtifact(artifact.absolutePath); ctx.ui.notify(`agy-artifacts: opened ${artifact.name}`, "info"); } catch (error) { ctx.ui.notify( `agy-artifacts: failed to open ${artifact.name} (${error instanceof Error ? error.message : error}).`, "error", ); } return; } if (arg) { ctx.ui.notify( 'agy-artifacts: usage "/agy-artifacts" or "/agy-artifacts open ".', "error", ); return; } await openAgyArtifactsPicker(ctx, rescan); }, }); pi.registerCommand("agy-subagents", { description: "List subagent activity observed on the agy stream (spawns, messages, kills)", handler: async (args, ctx) => { if (args.trim()) { ctx.ui.notify('agy-subagents: usage "/agy-subagents".', "error"); return; } const report = formatAgySubagents(subagentRoster); ctx.ui.notify( report ?? "antigravity: no subagent activity observed this session. Subagents appear here as agy streams invoke_subagent/send_message steps.", "info", ); }, }); pi.registerCommand("agy-usage", { description: "Show Antigravity model quotas (weekly and 5-hour limits)", handler: async (args, ctx) => { if (args.trim()) { ctx.ui.notify('agy-usage: usage "/agy-usage".', "error"); return; } await openAgyUsagePicker(ctx, (signal) => fetchAgyUsage({ signal })); }, }); }