/** * The lvrged_factory_* tool surface — eight resource-oriented tools: * * setup — toolchain detection + provider registry * status — one-screen infrastructure readout * ensure — idempotent "make model X runnable" entry point * pod — machine lifecycle: provision | register | pause | resume | destroy * job — generation lifecycle: run | progress | watch | finish * ledger — job queries + cost rollups * manifest — model/workflow registry entries * policy — the spend trust boundary (deliberately its own tool) * * Tools are deliberately thin: they own the state machine (registry, ledger, * idempotency) and the provider command templates, and they enforce the spend * policy. All *how-to* knowledge — installing runtimes, downloading weights, * diagnosing CUDA failures — lives in the skills, not here. */ import { Type } from "typebox"; import { execFile } from "node:child_process"; import type { ExtensionAPI } from "@earendil-works/pi-coding-agent"; import { loadProviders, saveProviders, loadMachines, saveMachines, loadDeployments, saveDeployments, loadModels, loadWorkflows, appendJob, loadJobs, loadPolicy, savePolicy, monthSpendUsd, todaySpendUsd, nextJobId, nextDeploymentId, nextMachineId, fmtUsd, stateDir, ProviderState, Deployment, Job, ModelManifest, WorkflowManifest, } from "./registry.js"; import { PROVIDERS, runSpec, cliExists, parseInstanceId, classifyCreateFailure, runpodGpuId, RUNPOD_DCS, H3_IMAGE, H3_CONTAINER_DISK_GB, } from "./providers.js"; type Ctx = any; // ExtensionCommandContext — kept loose to stay robust across pi versions function toolText(text: string, details: Record = {}) { return { content: [{ type: "text" as const, text }], details }; } /** Session-scoped call-to-call state for progress rate + stuck detection (no daemon). */ const progressState = new Map(); /** Background watchers: deployment/endpoint → interval + transition flags. */ interface WatchEntry { timer: ReturnType; base: string; promptId?: string; jobId?: string; intervalMs: number; label: string; lastLines: string[]; notifiedDone: boolean; notifiedUnreachable: boolean; notifiedStuck: boolean; } const watchers = new Map(); /** GET a ComfyUI JSON endpoint. http(s):// direct, or ssh://user@host:port for pods without a public HTTP port (tunnels curl over ssh — needs the runpodctl ssh key, override with LVRGED_FACTORY_SSH_KEY). */ async function comfyGet(url: string): Promise<{ ok: boolean; json: any }> { if (url.startsWith("ssh://")) { let u: URL; try { u = new URL(url); } catch { return { ok: false, json: null }; } const key = process.env.LVRGED_FACTORY_SSH_KEY || `${process.env.HOME || "/root"}/.runpod/ssh/runpodctl-ssh-key`; const path = `${u.pathname}${u.search}`; const cmd = `curl -s --max-time 6 "http://127.0.0.1:8188${path}"`; return new Promise((resolve) => { execFile("ssh", ["-o", "StrictHostKeyChecking=accept-new", "-o", "ConnectTimeout=6", "-o", "BatchMode=yes", "-i", key, "-p", u.port || "22", `${u.username || "root"}@${u.hostname}`, cmd], { timeout: 15000 }, (err, stdout) => { if (err || !stdout.trim()) return resolve({ ok: false, json: null }); try { resolve({ ok: true, json: JSON.parse(stdout) }); } catch { resolve({ ok: false, json: null }); } }); }); } try { const r = await fetch(url, { signal: AbortSignal.timeout(8000) }); if (!r.ok) return { ok: false, json: null }; return { ok: true, json: await r.json() }; } catch { return { ok: false, json: null }; } } /** One poll of a ComfyUI endpoint → phase/percent/ETA/queue/stuck readout. Shared by job action=progress and action=watch. */ async function pollComfyProgress(base: string, promptId: string | undefined, job: Job | undefined, costHr: number): Promise<{ lines: string[]; phase: string; pct: number; eta_s: number | null; queueRemaining: number; stuck: boolean; unreachable: boolean; done: boolean; cost: number; }> { let done = false; if (promptId) { const h = await comfyGet(`${base}/history/${promptId}`); if (h.ok) { const hist = h.json as Record | undefined; const entry = hist?.[promptId]; done = !!entry?.outputs && Object.keys(entry.outputs).length > 0; } } const pr = await comfyGet(`${base}/progress`); const queue = await comfyGet(`${base}/queue`); const unreachable = !pr.ok; // ComfyUI /progress returns a 0..1 fraction of the current item's sampling steps const frac = typeof pr.json?.progress === "number" ? Math.min(1, Math.max(0, pr.json.progress)) : 0; const queueRemaining = typeof queue.json?.queue_pending?.length === "number" ? queue.json.queue_pending.length : typeof pr.json?.queue_remaining === "number" ? pr.json.queue_remaining : 0; const key = `${base}:${promptId || ""}`; const prev = progressState.get(key); const now = Date.now(); const elapsed = job ? (now - new Date(job.started_at).getTime()) / 1000 : 0; let phase: string; if (unreachable) phase = "unreachable"; else if (done) phase = "done"; else if (frac >= 1) phase = "decode"; else if (frac > 0) phase = "sample"; else phase = "load"; const rate = prev && now > prev.ts ? (frac - prev.frac) / ((now - prev.ts) / 1000) : 0; let eta_s: number | null = null; if (phase === "sample" && frac > 0 && rate > 0) eta_s = Math.round((1 - frac) / rate); const stuck = !!prev && phase === prev.phase && Math.abs(frac - prev.frac) < 0.0001 && now - prev.ts > 180000; progressState.set(key, { ts: now, frac, phase, elapsed }); const cost = job && elapsed > 0 ? Math.round((elapsed / 3600) * costHr * 1000) / 1000 : 0; const bar = "▓".repeat(Math.round(frac * 10)) + "░".repeat(10 - Math.round(frac * 10)); const lines = [ `job ${job?.id || "?"} · ${phase} ${Math.round(frac * 100)}% ${bar}${eta_s ? ` · eta ~${fmtDur(eta_s)}` : ""}`, `${cost > 0 ? ` cost so far ${fmtUsd(cost)}` : ""}${queueRemaining > 0 ? ` · ${queueRemaining} ahead in queue` : ""}${stuck ? " · ⚠ stuck 3+ min" : ""}`.trim(), ].filter(Boolean); return { lines, phase, pct: Math.round(frac * 100), eta_s, queueRemaining, stuck, unreachable, done, cost }; } function fmtDur(s: number): string { if (s < 60) return `${Math.max(1, Math.round(s))}s`; const m = Math.floor(s / 60); return `${m}m ${Math.round(s % 60)}s`; } /** Ask the user before anything that spends real money above the policy threshold. */ async function spendGate(ctx: Ctx, what: string, estUsd: number): Promise { const policy = loadPolicy(); if (estUsd <= policy.confirm_above_usd) return true; const ok = await ctx.ui.confirm("LVRGED Factory GPU spend", `${what} is estimated at ${fmtUsd(estUsd)}. Proceed?`, { detail: `Policy: confirm above ${fmtUsd(policy.confirm_above_usd)}.`, acceptLabel: "Proceed", rejectLabel: "Cancel", }); return ok; } export function registerFactoryTools(pi: ExtensionAPI) { // ------------------------------------------------------------------ setup pi.registerTool({ name: "lvrged_factory_setup", label: "LVRGED Factory Setup", description: "Detect the required toolchain — runpodctl (RunPod provider) and hf (Hugging Face CLI for weight downloads) — check authentication, and write the provider registry. Run this once per machine; re-run to rescan. If a CLI is missing, the skill lvrged-factory-provider-adapters tells you how to install it (ask permission first).", parameters: Type.Object({ rescan: Type.Optional(Type.Boolean({ description: "Force re-detection of CLIs and auth (default false — uses cached results)" })), }), async execute(_id, params: { rescan?: boolean }, _sig, _upd, ctx: Ctx) { const providers = loadProviders(); const out: string[] = ["LVRGED-FACTORY SETUP", "-------------------"]; const updated: ProviderState[] = []; const now = new Date().toISOString(); for (const [pid, man] of Object.entries(PROVIDERS)) { const cached = providers.find((p) => p.id === pid); let installed = cached?.installed ?? false; let version: string | undefined = cached?.cliVersion; if (!cached || params.rescan) { if (man.cli) installed = await cliExists(man.cli); } else { installed = cached.installed; } let authenticated = cached?.authenticated ?? false; let authHint: string | undefined; if (installed && man.authCheck) { const r = await runSpec(man.authCheck, {}, 15000); authenticated = r.ok; if (!r.ok) authHint = (r.stderr || r.stdout || r.error || "").slice(0, 120); } else if (installed) { authenticated = true; // no auth check defined — treat CLI presence as enough } updated.push({ id: pid, name: man.name, kind: man.kind, cli: man.cli, installed, cliVersion: version, authenticated, authHint, detectedAt: now }); const mark = installed ? "✓" : "○"; const auth = authenticated ? "authenticated" : "NOT authenticated"; out.push(`${mark} ${man.name.padEnd(10)} ${installed ? (man.cli || "").padEnd(10) + " " + auth : "CLI not installed"}`); if (authHint) out.push(` hint: ${authHint}`); } saveProviders(updated); const ready = updated.filter((p) => p.installed && p.authenticated); // Hugging Face is part of the toolchain (weight downloads) but not a provider const hfInstalled = await cliExists("hf"); let hfAuth = false; let hfHint: string | undefined; if (hfInstalled) { const r = await runSpec({ cmd: "hf", args: () => ["auth", "whoami"] }, {}, 15000); hfAuth = r.ok; if (!r.ok) hfHint = (r.stderr || r.stdout || r.error || "").slice(0, 120); } out.push(`${hfInstalled ? "✓" : "○"} ${hfInstalled ? "Hugging Face".padEnd(10) + "hf".padEnd(10) + (hfAuth ? "authenticated" : "NOT authenticated") : "Hugging Face hf not installed (needed for weight downloads)"}`); if (hfHint) out.push(` hint: ${hfHint}`); out.push("-----------------"); out.push(`Available providers: ${ready.length ? ready.map((p) => p.name).join(", ") : "none yet"}`); out.push("Install missing CLIs via the lvrged-factory-provider-adapters skill, then re-run /lvrged-factory setup."); return toolText(out.join("\n"), { providers: updated }); }, }); // ----------------------------------------------------------------- status pi.registerTool({ name: "lvrged_factory_status", label: "LVRGED Factory Status", description: "Report the current GPU infrastructure: providers, deployments (with status), queued/running/completed jobs, and spend today and this month. Use before deciding to provision, reuse, pause, or destroy anything.", parameters: Type.Object({}), async execute() { const providers = loadProviders(); const deployments = loadDeployments(); const jobs = loadJobs(); const out: string[] = ["LVRGED-FACTORY GPU", "------------------"]; out.push("Providers"); for (const p of providers) out.push(` ${p.name.padEnd(10)} ${p.installed && p.authenticated ? "✓ ready" : p.installed ? "○ installed, not authed" : "○ not installed"}`); out.push("Deployments"); if (!deployments.length) out.push(" (none)"); for (const d of deployments) { const dot = d.status === "ready" ? "●" : d.status === "provisioning" ? "◐" : d.status === "stopped" ? "○" : "✕"; out.push(` ${dot} ${d.id.padEnd(14)} ${d.model || d.runtime} on ${d.gpu} @ ${d.provider} ${d.endpoint ? d.endpoint : ""} (${fmtUsd(d.cost_hr)}/hr)`); } out.push(`Jobs ${jobs.filter((j) => j.status === "running").length} running · ${jobs.filter((j) => j.status === "queued").length} queued · ${jobs.filter((j) => j.status === "completed").length} completed`); out.push(`Spend Today ${fmtUsd(todaySpendUsd())} · Month ${fmtUsd(monthSpendUsd())}`); const policy = loadPolicy(); out.push(`Policy: per-job ≤ ${fmtUsd(policy.ceiling_per_job_usd)} · daily ≤ ${fmtUsd(policy.ceiling_daily_usd)} · monthly ≤ ${fmtUsd(policy.ceiling_monthly_usd)}`); return toolText(out.join("\n"), { deployments, providers }); }, }); // ----------------------------------------------------------------- ensure pi.registerTool({ name: "lvrged_factory_ensure", label: "LVRGED Factory Ensure", description: "The idempotent entry point: given a model (and optionally workflow constraints), return a READY deployment that can run it — reusing an existing one, surfacing a paused pod worth resuming, or producing a provisioning plan. Prefer this over lvrged_factory_pod action=provision when the goal is 'make model X runnable'.", parameters: Type.Object({ model: Type.String({ description: "Model id from models.json (e.g. wan-2.2-5b, minimax-h3, hunyuan-video)" }), workflow: Type.Optional(Type.String({ description: "Workflow id to run (e.g. h3-image-to-video)" })), gpu: Type.Optional(Type.String({ description: "Force a GPU class (default: cheapest class that fits the model's VRAM)" })), provider: Type.Optional(Type.String({ description: "Force a provider" })), cost_max_hr: Type.Optional(Type.Number({ description: "Max hourly rate (USD)" })), create: Type.Optional(Type.Boolean({ description: "Produce a provisioning plan if nothing ready exists (default true)" })), }), async execute(_id, p: any) { const models = loadModels(); const model = models.find((m) => m.id === p.model); if (!model) { return toolText(`Model "${p.model}" not in registry. Add it with lvrged_factory_manifest kind=model (or read the lvrged-factory-model-deployment skill to discover it first).`); } const deployments = loadDeployments(); const ready = deployments.filter((d) => d.status === "ready" && (!p.workflow || d.model === p.model) && (!p.gpu || d.gpu === p.gpu)); const reuse = ready.find((d) => !p.provider || d.provider === p.provider) || ready[0]; if (reuse) { return toolText( [ `REUSE ${reuse.id}`, ` ${reuse.gpu} on ${reuse.provider} — ${reuse.runtime}${reuse.model ? ` for ${reuse.model}` : ""}`, ` endpoint: ${reuse.endpoint || "—"}`, ` cost: ${fmtUsd(reuse.cost_hr)}/hr`, ].join("\n"), { deployment: reuse, action: "reuse" }, ); } if (p.create === false) return toolText(`No ready deployment for ${p.model}.`, { action: "none" }); const fits = model.vram_gb_fp8 || model.vram_gb_min; const gpuSuggestion = p.gpu || defaultGpuForVram(fits); const providerSuggestion = p.provider || cheapestProviderFor(gpuSuggestion); // paused pod first — resuming beats reprovisioning when disk state survives const paused = deployments.find((d) => d.status === "stopped" && d.model === p.model); const plan = [ `PLAN: provision ${model.name} (${model.gpu_class}, ≥${model.vram_gb_min}GB VRAM)`, paused ? ` 0. Paused pod exists: ${paused.id} (${paused.instance_id}) — lvrged_factory_pod action=resume it first; cheaper than a fresh pod if the disk survived.` : "", ` 1. lvrged_factory_pod action=provision provider=${providerSuggestion} gpu=${gpuSuggestion} — defaults to the cu130 H3 image + capacity fallback; just call it`, ` 2. Install runtime + weights per the deployment/comfyui skills (agent-driven)`, ` 3. lvrged_factory_pod action=register with endpoint/ssh → becomes READY`, ` 4. lvrged_factory_job action=run with workflow ${p.workflow || ""}`, ` Reference pricing: docs/pricing/${providerSuggestion}.md (snapshot — verify live before committing spend)`, ].filter(Boolean); return toolText(plan.join("\n"), { action: "provision", plan: { provider: providerSuggestion, gpu: gpuSuggestion } }); }, }); // -------------------------------------------------------------------- pod pi.registerTool({ name: "lvrged_factory_pod", label: "LVRGED Factory Pod", description: "The machine lifecycle, one tool: action=provision (capacity-fallback create: COMMUNITY → SECURE, DC rotation, disk shrink; defaults to the H3 lane — RTX PRO 6000, cu130 image, 80GB disk, ports 22+8188), action=register (record a ready instance as the stable deployment handle; idempotent), action=pause (pod stop: GPU billing off, pod kept; EXITED pods vanish from pod list), action=resume (pod start: billing restarts, SSH port usually changes, verify weights survived), action=destroy (terminate; ledger history kept). Money actions confirm with the user per the spend policy.", parameters: Type.Object({ action: Type.String({ description: "provision | register | pause | resume | destroy", enum: ["provision", "register", "pause", "resume", "destroy"] }), // shared provider: Type.Optional(Type.String({ description: "Provider id (default runpod)" })), id: Type.Optional(Type.String({ description: "pause/resume/destroy: deployment id (dep_...) or provider instance id" })), // provision gpu: Type.Optional(Type.String({ description: "provision/register: GPU friendly name (default RTX PRO 6000); mapped to the exact RunPod --gpu-id string" })), model: Type.Optional(Type.String({ description: "provision/register: model id this instance is for (default minimax-h3)" })), image: Type.Optional(Type.String({ description: "provision: docker image (default: the cu130 H3 image). Overrides template." })), template: Type.Optional(Type.String({ description: "provision: RunPod template id — only if you explicitly don't want the H3 image (cu128 templates are ~2x slower for H3 INT8)" })), cloud: Type.Optional(Type.String({ description: "provision: starting pool — COMMUNITY (default; falls back to SECURE automatically) or SECURE (skip community)", enum: ["COMMUNITY", "SECURE"] })), disk_gb: Type.Optional(Type.Number({ description: "provision: container disk GB (default 80; no network volume, ever)" })), dc: Type.Optional(Type.String({ description: "provision: pin one datacenter id (skips rotation), e.g. US-NC-1" })), cost_max_hr: Type.Optional(Type.Number({ description: "provision: refuse above this hourly rate (USD)" })), // register instance_id: Type.Optional(Type.String({ description: "register: provider's instance id (required if machine_id not given)" })), machine_id: Type.Optional(Type.String({ description: "register: machine id from a previous provision" })), runtime: Type.Optional(Type.String({ description: "register/provision: runtime — comfyui (default) | native | custom" })), endpoint: Type.Optional(Type.String({ description: "register: HTTP endpoint (ComfyUI: http://127.0.0.1:8188 via tunnel, or ssh://user@host:port)" })), ssh: Type.Optional(Type.String({ description: "register: SSH target, e.g. root@1.2.3.4" })), cost_hr: Type.Optional(Type.Number({ description: "register: actual hourly cost (USD)" })), status: Type.Optional(Type.String({ description: "register: ready (default) | provisioning | stopped | error" })), }), async execute(_id, p: any, _sig, _upd, ctx: Ctx) { const providerId = p.provider || "runpod"; if (p.action === "provision") { const man = PROVIDERS[providerId]; if (!man) return toolText(`Unknown provider: ${providerId}`); if (!man.provision) return toolText(`${man.name}: no CLI provisioning template. See lvrged-factory-provider-adapters skill.`); const gpuFriendly = p.gpu || "RTX PRO 6000"; const gpuId = providerId === "runpod" ? runpodGpuId(gpuFriendly) : gpuFriendly; const clouds = p.cloud === "SECURE" ? ["SECURE"] : ["COMMUNITY", "SECURE"]; const worstHr = estimateHourly(providerId, gpuFriendly, "SECURE"); if (p.cost_max_hr && worstHr > p.cost_max_hr) { return toolText(`Refusing: ${gpuFriendly} on ${man.name} ≈ ${fmtUsd(worstHr)}/hr (secure) exceeds cost_max_hr ${fmtUsd(p.cost_max_hr)}. Pick a cheaper GPU or raise the cap.`); } if (!(await spendGate(ctx, `Provisioning ${gpuFriendly} on ${man.name} (up to ${fmtUsd(worstHr)}/hr if community is out of stock)`, worstHr * 2))) { return toolText("Cancelled by user."); } const image = p.template ? "" : (p.image || H3_IMAGE); const dcs = p.dc ? [p.dc] : RUNPOD_DCS; let disk = p.disk_gb || H3_CONTAINER_DISK_GB; const attempts: string[] = []; let created: { ok: boolean; stdout: string; cloud: string; dc: string } | null = null; outer: for (const cloud of clouds) { let shrunk = false; for (const dc of dcs) { const vars: Record = { gpu: gpuId, cloud, dc, disk: String(disk), image, template: p.template || "", name: `lvrged-${(p.model || "gpu").replace(/[^a-z0-9]/gi, "")}`, }; const r = await runSpec(man.provision, vars, 120000); const out = r.stdout + r.stderr + (r.error || ""); if (r.ok && parseInstanceId(providerId, r.stdout)) { created = { ok: true, stdout: r.stdout, cloud, dc }; break outer; } const kind = classifyCreateFailure(out); attempts.push(`${cloud}/${dc}/disk${disk}: ${kind === "no_stock" ? "no stock" : kind === "spec_too_big" ? "spec too big" : out.slice(0, 90)}`); if (kind === "spec_too_big" && !shrunk && disk > 60) { // machine pool exists but the disk doesn't fit — shrink once and retry this DC disk = 60; shrunk = true; const r2 = await runSpec(man.provision, { ...vars, disk: String(disk) }, 120000); if (r2.ok && parseInstanceId(providerId, r2.stdout)) { created = { ok: true, stdout: r2.stdout, cloud, dc }; break outer; } attempts.push(`${cloud}/${dc}/disk${disk}: ${classifyCreateFailure(r2.stdout + r2.stderr) === "no_stock" ? "no stock" : (r2.stderr || r2.stdout).slice(0, 90)}`); } if (kind === "other") break; // auth/flag error — rotating DCs won't help } } if (!created) { return toolText( [ `PROVISION FAILED — capacity protocol exhausted (${attempts.length} attempts)`, ...attempts.slice(0, 12).map((a) => ` ${a}`), ` Nothing was created; failed creates don't bill.`, ` Next: wait 10-30 min and retry, try another GPU class, or check the pool by hand — see the lvrged-factory-provider-adapters skill.`, ].join("\n"), { attempts }, ); } const instanceId = parseInstanceId(providerId, created.stdout)!; const actualHr = estimateHourly(providerId, gpuFriendly, created.cloud); const machineId = nextMachineId(); const machines = loadMachines(); machines.push({ id: machineId, provider: providerId, instance_id: instanceId, gpu: gpuFriendly, gpu_count: 1, status: "running", cost_hr: actualHr, region: created.dc, created_at: new Date().toISOString(), tags: [p.model || p.runtime || "comfyui"], }); saveMachines(machines); const depId = nextDeploymentId(p.model || "minimax-h3"); const deployments = loadDeployments(); deployments.push({ id: depId, provider: providerId, machine_id: machineId, instance_id: instanceId, gpu: gpuFriendly, gpu_count: 1, runtime: p.runtime || "comfyui", model: p.model || "minimax-h3", status: "provisioning", cost_hr: actualHr, region: created.dc, created_at: new Date().toISOString(), }); saveDeployments(deployments); const lines = [ `PROVISION ${depId} — pod ${instanceId} (${created.cloud} @ ${created.dc}, ${fmtUsd(actualHr)}/hr, billing started)`, ` gpu: ${gpuFriendly} (${gpuId})`, ` image: ${image || `template ${p.template}`} · disk ${disk}GB · ports 22/tcp,8188/http`, attempts.length ? ` (${attempts.length} capacity rejections before this one — normal)` : "", ` Next (comfyui skill has the exact commands):`, ` 1. Poll \`runpodctl pod get ${instanceId}\` until .ssh.ssh_command appears — the image pull takes minutes; "Connection refused" while it extracts is normal.`, ` 2. Install weights + nodes, put flags in the template's comfyui_args.txt, RESTART THE POD to apply — never hand-roll nohup daemons over SSH.`, ` 3. Port 8188 often isn't publicly reachable — use the SSH tunnel (ssh -f -N -L 8188:127.0.0.1:8188 ...) and endpoint http://127.0.0.1:8188.`, ` 4. Health-check with a test generation, then lvrged_factory_pod action=register.`, ].filter(Boolean); return toolText(lines.join("\n"), { deployment_id: depId, machine_id: machineId, instance_id: instanceId, cloud: created.cloud, dc: created.dc, disk_gb: disk, attempts }); } if (p.action === "register") { if (!p.gpu || !p.runtime) return toolText("register needs gpu and runtime (plus instance_id or machine_id)."); const deployments = loadDeployments(); const existing = deployments.find((d) => d.instance_id === p.instance_id && d.instance_id !== "pending") || deployments.find((d) => d.machine_id === p.machine_id); let dep: Deployment; if (existing) { dep = { ...existing, gpu: p.gpu || existing.gpu, runtime: p.runtime || existing.runtime, model: p.model || existing.model, endpoint: p.endpoint || existing.endpoint, ssh: p.ssh || existing.ssh, cost_hr: p.cost_hr ?? existing.cost_hr, status: (p.status || "ready") as Deployment["status"] }; saveDeployments(deployments.map((d) => (d.id === existing.id ? dep : d))); } else { dep = { id: nextDeploymentId(p.model), provider: providerId, instance_id: p.instance_id, gpu: p.gpu, gpu_count: 1, runtime: p.runtime, model: p.model, endpoint: p.endpoint, ssh: p.ssh, status: (p.status || "ready") as Deployment["status"], cost_hr: p.cost_hr ?? estimateHourly(providerId, p.gpu), created_at: new Date().toISOString(), }; deployments.push(dep); saveDeployments(deployments); } const lines = [ `DEPLOYMENT ${dep.id}`, ` provider: ${dep.provider}`, ` gpu: ${dep.gpu}`, ` runtime: ${dep.runtime}`, ` model: ${dep.model || "—"}`, ` endpoint: ${dep.endpoint || "—"}`, ` ssh: ${dep.ssh || "—"}`, ` cost: ${fmtUsd(dep.cost_hr)}/hr`, ` status: ${dep.status.toUpperCase()}`, ]; return toolText(lines.join("\n"), { deployment: dep }); } if (p.action === "pause") { if (!p.id) return toolText("pause needs id (deployment id or instance id)."); const deployments = loadDeployments(); const dep = deployments.find((d) => d.id === p.id) || deployments.find((d) => d.instance_id === p.id); const instanceId = dep?.instance_id || p.id; const man = PROVIDERS[dep?.provider || providerId]; if (!man?.stop) return toolText(`No stop command for provider ${dep?.provider || providerId}.`); const r = await runSpec(man.stop, { id: instanceId }, 60000); if (dep) saveDeployments(deployments.map((d) => (d.id === dep.id ? { ...d, status: "stopped" as const, endpoint: undefined, idle_since: new Date().toISOString() } : d))); return toolText( r.ok ? `PAUSED ${dep?.id || instanceId} — GPU billing stopped, pod kept (EXITED; gone from pod list, still in pod get).\nResume with lvrged_factory_pod action=resume; expect the SSH port to change and weights to possibly need re-download.` : `Pause failed: ${(r.error || r.stderr).slice(0, 200)}`, { ok: r.ok }, ); } if (p.action === "resume") { if (!p.id) return toolText("resume needs id (deployment id or instance id)."); const deployments = loadDeployments(); const dep = deployments.find((d) => d.id === p.id) || deployments.find((d) => d.instance_id === p.id); const instanceId = dep?.instance_id || p.id; const man = PROVIDERS[dep?.provider || providerId]; if (!man?.start) return toolText(`No start command for provider ${dep?.provider || providerId}.`); if (!(await spendGate(ctx, `Resuming pod ${instanceId} (${fmtUsd(dep?.cost_hr || 2)}/hr billing restarts)`, dep?.cost_hr || 2))) { return toolText("Cancelled by user."); } const r = await runSpec(man.start, { id: instanceId }, 60000); if (r.ok && dep) saveDeployments(deployments.map((d) => (d.id === dep.id ? { ...d, status: "provisioning" as const, idle_since: null } : d))); return toolText( r.ok ? `RESUMING ${dep?.id || instanceId} — billing restarted.\nNext: poll \`runpodctl pod get ${instanceId}\` for the new ssh_command, verify weights on disk, health-check, then lvrged_factory_pod action=register status=ready with the (possibly new) endpoint.` : `Resume failed: ${(r.error || r.stderr).slice(0, 200)}`, { ok: r.ok }, ); } if (p.action === "destroy") { if (!p.id) return toolText("destroy needs id (deployment id dep_... or machine id m_...)."); const deployments = loadDeployments(); const machines = loadMachines(); const dep = deployments.find((d) => d.id === p.id); const machine = machines.find((m) => m.id === p.id) || (dep ? machines.find((m) => m.id === dep.machine_id) : undefined); if (!dep && !machine) return toolText(`Nothing found for "${p.id}".`); const recentJobs = loadJobs().filter((j) => (dep && j.deployment_id === dep.id) && (j.status === "running" || j.status === "queued")); if (recentJobs.length) { const ok = await ctx.ui.confirm("LVRGED Factory GPU destroy", `${recentJobs.length} job(s) still running/queued on ${p.id}. Destroy anyway?`, { acceptLabel: "Destroy", rejectLabel: "Cancel" }); if (!ok) return toolText("Cancelled."); } const man = dep ? PROVIDERS[dep.provider] : undefined; let cliResult = ""; if (man?.destroy && machine?.instance_id && machine.instance_id !== "pending") { const r = await runSpec(man.destroy, { id: machine.instance_id }, 60000); cliResult = r.ok ? "destroy command ok" : `destroy command failed: ${(r.error || r.stderr).slice(0, 200)}`; } if (machine) saveMachines(machines.map((m) => (m.id === machine.id ? { ...m, status: "destroyed" as const } : m))); if (dep) saveDeployments(deployments.map((d) => (d.id === dep.id ? { ...d, status: "destroyed" as const, endpoint: undefined } : d))); return toolText(`Destroyed ${p.id}. ${cliResult}\nRecord kept (status destroyed). Jobs and cost history remain in the ledger.`); } return toolText(`Unknown action "${p.action}". Use provision | register | pause | resume | destroy.`); }, }); // -------------------------------------------------------------------- job pi.registerTool({ name: "lvrged_factory_job", label: "LVRGED Factory Job", description: "The generation lifecycle, one tool: action=run (open the ledger entry, then execute per the workflow skill: POST /prompt, capture prompt_id), action=progress (one-shot poll: phase load/sample/decode/done/stuck, percent, ETA, queue position, cost so far; call every 5-10s), action=watch (background queue watcher: polls /progress + /queue on its own, keeps the TUI widget live, notifies on done/stuck/unreachable — start it, walk away; sub-actions watch_stop / watch_status), action=finish (close the ledger: status, duration, cost, artifact URIs, and the observed load/encode/sample/decode phase timings that calibrate the ETA model). Endpoints may be http(s):// or ssh://user@host:port.", parameters: Type.Object({ action: Type.String({ description: "run | progress | watch | watch_stop | watch_status | finish", enum: ["run", "progress", "watch", "watch_stop", "watch_status", "finish"] }), // shared deployment_id: Type.Optional(Type.String({ description: "Deployment id (dep_...) — required for run, default source of the endpoint for progress/watch" })), job_id: Type.Optional(Type.String({ description: "Job id (J-...) — required for finish; enables cost-so-far for progress/watch" })), prompt_id: Type.Optional(Type.String({ description: "ComfyUI prompt_id from POST /prompt — enables done detection" })), endpoint: Type.Optional(Type.String({ description: "progress/watch: endpoint override, e.g. http://127.0.0.1:8188 or ssh://root@1.2.3.4:22022" })), // run workflow: Type.Optional(Type.String({ description: "run: workflow id (e.g. h3-text-to-video)" })), inputs: Type.Optional(Type.Object({}, { additionalProperties: true, description: "run: workflow inputs — prompt, image, duration, etc." })), model: Type.Optional(Type.String({ description: "run: model id (defaults to the deployment's model)" })), artifacts_expected: Type.Optional(Type.Array(Type.String({ description: "run: expected output names, e.g. output.mp4" }))), // watch interval_s: Type.Optional(Type.Number({ description: "watch: poll interval seconds (default 8, min 3)" })), // finish status: Type.Optional(Type.String({ description: "finish: completed | failed | cancelled", enum: ["completed", "failed", "cancelled"] })), artifacts: Type.Optional(Type.Array(Type.Object({ name: Type.String(), uri: Type.String() }), { description: "finish: output artifacts (durable URIs — copy off the pod before shutdown!)" })), duration_s: Type.Optional(Type.Number({ description: "finish: wall-clock duration in seconds" })), cost_usd: Type.Optional(Type.Number({ description: "finish: actual cost; omit to estimate from deployment rate" })), error: Type.Optional(Type.String({ description: "finish: failure reason" })), workflow_version: Type.Optional(Type.Number({ description: "finish: workflow version used" })), load_s: Type.Optional(Type.Number({ description: "finish: seconds from POST to first progress tick (weight load)" })), encode_s: Type.Optional(Type.Number({ description: "finish: seconds of encode before sampling" })), sample_s: Type.Optional(Type.Number({ description: "finish: seconds of sampling (steps)" })), decode_s: Type.Optional(Type.Number({ description: "finish: seconds from last tick to output file (decode + assemble)" })), }), async execute(_id, p: any, _sig, _upd, ctx: Ctx) { if (p.action === "run") { if (!p.deployment_id || !p.workflow) return toolText("run needs deployment_id and workflow."); const deployments = loadDeployments(); const dep = deployments.find((d) => d.id === p.deployment_id); if (!dep) return toolText(`Unknown deployment ${p.deployment_id}.`); if (dep.status !== "ready") return toolText(`Deployment ${dep.id} is ${dep.status}, not ready. Ensure it first.`); const now = new Date().toISOString(); const job: Job = { id: nextJobId(), deployment_id: dep.id, provider: dep.provider, model: p.model || dep.model, workflow: p.workflow, inputs: p.inputs || {}, status: "running", started_at: now, gpu: dep.gpu, instance_id: dep.instance_id, artifacts: [], ts: now, }; appendJob(job); try { ctx.ui.setWidget("lvrged-factory-gpu-progress", [ `job ${job.id} · ${dep.gpu} @ ${dep.provider}`, ` queued — POST the workflow, capture prompt_id, then lvrged_factory_job action=watch`, ]); } catch { /* widget unavailable — non-interactive mode */ } return toolText( [ `JOB ${job.id}`, ` workflow: ${job.workflow}`, ` deployment: ${dep.id} (${dep.gpu} @ ${dep.provider})`, ` started: ${now}`, ` inputs: ${JSON.stringify(job.inputs).slice(0, 200)}`, `Execute per the workflow skill, then lvrged_factory_job action=finish job_id=${job.id} with artifact URIs.`, ].join("\n"), { job }, ); } if (p.action === "progress") { const deployments = loadDeployments(); const dep = p.deployment_id ? deployments.find((d) => d.id === p.deployment_id) : undefined; if (!dep && !p.endpoint) return toolText("progress needs deployment_id or endpoint."); const base = (p.endpoint || dep?.endpoint || "").replace(/\/$/, ""); if (!base) return toolText(`Deployment ${dep?.id} has no endpoint — pass endpoint explicitly or register it.`); const job = p.job_id ? loadJobs().find((j) => j.id === p.job_id) : undefined; const r = await pollComfyProgress(base, p.prompt_id, job, dep?.cost_hr || 0); const lines = [ `job ${job?.id || "?"} · ${dep ? `${dep.gpu} @ ${dep.provider}` : base}${r.done ? " · ✓ done" : ""}`, ...r.lines, r.unreachable ? " ⚠ endpoint unreachable — is ComfyUI up? (ssh:// endpoint or tunnel needed)" : "", ].filter(Boolean); try { if (r.done || r.unreachable) ctx.ui.setWidget("lvrged-factory-gpu-progress", undefined); else ctx.ui.setWidget("lvrged-factory-gpu-progress", lines); } catch { /* widget unavailable */ } return toolText(lines.join("\n"), { phase: r.phase, pct: r.pct, eta_s: r.eta_s, queue_remaining: r.queueRemaining, stuck: r.stuck, cost_so_far_usd: r.cost, unreachable: r.unreachable, done: r.done }); } if (p.action === "watch_stop") { const key = p.deployment_id || p.endpoint; const w = key ? watchers.get(key) : undefined; if (!w) return toolText(`No watcher on ${key || "?"}.`); clearInterval(w.timer); watchers.delete(key); try { ctx.ui.setWidget("lvrged-factory-gpu-progress", undefined); } catch { /* no UI */ } return toolText(`Stopped watching ${w.label}.`); } if (p.action === "watch_status") { if (!watchers.size) return toolText("No background watchers running. Start one with action=watch."); const lines = [...watchers.entries()].map(([k, w]) => ` ${k} — ${w.lastLines.join(" | ") || "waiting for first poll"}`); return toolText(["WATCHERS", ...lines].join("\n"), { watchers: [...watchers.keys()] }); } if (p.action === "watch") { const deployments = loadDeployments(); const dep = p.deployment_id ? deployments.find((d) => d.id === p.deployment_id) : undefined; if (!dep && !p.endpoint) return toolText("Give deployment_id (registry) or endpoint (http:// or ssh://) to watch."); const base = (p.endpoint || dep?.endpoint || "").replace(/\/$/, ""); if (!base) return toolText(`Deployment ${dep?.id} has no endpoint — pass endpoint explicitly.`); const intervalMs = Math.max(3000, (p.interval_s || 8) * 1000); const key = p.deployment_id || base; const label = dep ? `${dep.id} (${dep.gpu} @ ${dep.provider})` : base; const old = watchers.get(key); if (old) clearInterval(old.timer); const entry: WatchEntry = { timer: null as any, base, promptId: p.prompt_id, jobId: p.job_id, intervalMs, label, lastLines: [], notifiedDone: false, notifiedUnreachable: false, notifiedStuck: false }; const tick = async () => { const j = entry.jobId ? loadJobs().find((x) => x.id === entry.jobId) : undefined; const r = await pollComfyProgress(entry.base, entry.promptId, j, dep?.cost_hr || 0); entry.lastLines = r.lines; try { if (r.done) ctx.ui.setWidget("lvrged-factory-gpu-progress", [`${entry.label} · ✓ done`, ...r.lines]); else if (r.unreachable) ctx.ui.setWidget("lvrged-factory-gpu-progress", [`${entry.label}`, " ⚠ unreachable — is ComfyUI up yet?", ...r.lines]); else ctx.ui.setWidget("lvrged-factory-gpu-progress", [`${entry.label}`, ...r.lines]); } catch { /* widget unavailable */ } if (r.done && !entry.notifiedDone) { entry.notifiedDone = true; try { ctx.ui.notify(`lvrged-factory: ${entry.label} — generation done ✓`, "info"); } catch { /* no UI */ } clearInterval(entry.timer); watchers.delete(key); } if (r.unreachable && !entry.notifiedUnreachable) { entry.notifiedUnreachable = true; try { ctx.ui.notify(`lvrged-factory: ${entry.label} unreachable — retrying in the background`, "warning"); } catch { /* no UI */ } } if (!r.unreachable) entry.notifiedUnreachable = false; if (r.stuck && !entry.notifiedStuck) { entry.notifiedStuck = true; try { ctx.ui.notify(`lvrged-factory: ${entry.label} stuck — no progress for 3+ min`, "warning"); } catch { /* no UI */ } } if (!r.stuck) entry.notifiedStuck = false; }; entry.timer = setInterval(tick, intervalMs); watchers.set(key, entry); void tick(); return toolText( `Watching ${label} every ${intervalMs / 1000}s in the background — the widget stays live and you'll be notified on done. Stop with lvrged_factory_job action=watch_stop${p.job_id ? " job_id=" + p.job_id : ""}.`, { watching: key }, ); } if (p.action === "finish") { if (!p.job_id || !p.status) return toolText("finish needs job_id and status (completed | failed | cancelled)."); const jobs = loadJobs(); const idx = jobs.findIndex((j) => j.id === p.job_id); if (idx < 0) return toolText(`Unknown job ${p.job_id}.`); const job = jobs[idx]; const dep = loadDeployments().find((d) => d.id === job.deployment_id); const duration = p.duration_s ?? ((Date.now() - new Date(job.started_at).getTime()) / 1000); const cost = p.cost_usd ?? (p.status === "completed" ? Math.max(0.001, (duration / 3600) * (dep?.cost_hr || 1)) : 0); const completed: Job = { ...job, status: p.status, completed_at: new Date().toISOString(), duration_s: Math.round(duration * 10) / 10, cost_usd: Math.round(cost * 1000) / 1000, artifacts: p.artifacts || job.artifacts, error: p.error, workflow_version: p.workflow_version ?? job.workflow_version, prompt_id: p.prompt_id ?? job.prompt_id, load_s: p.load_s, encode_s: p.encode_s, sample_s: p.sample_s, decode_s: p.decode_s, }; try { const done = p.status === "completed" ? "✓ done" : p.status === "failed" ? "✕ failed" : "—"; ctx.ui.setWidget("lvrged-factory-gpu-progress", [ `job ${job.id} ${done} · ${(completed.duration_s ?? 0).toFixed(0)}s · ${fmtUsd(completed.cost_usd || 0)}`, completed.artifacts.length ? ` artifacts: ${completed.artifacts.map((a) => a.uri).join(", ")}` : ` error: ${p.error || "none recorded"}`, completed.sample_s ? ` phases: load ${completed.load_s ?? "?"}s · encode ${completed.encode_s ?? "?"}s · sample ${completed.sample_s}s · decode ${completed.decode_s ?? "?"}s` : "", ].filter(Boolean)); } catch { /* widget unavailable */ } // rewrite the JSONL line in place (keeps the ledger single-source) const fs = await import("node:fs"); const path = `${stateDir()}/jobs.jsonl`; const lines = fs.readFileSync(path, "utf-8").split("\n").filter(Boolean); lines[idx] = JSON.stringify(completed); fs.writeFileSync(path, lines.join("\n") + "\n"); return toolText( [ `JOB ${job.id} → ${p.status.toUpperCase()}`, ` duration: ${completed.duration_s}s`, ` cost: ${fmtUsd(completed.cost_usd || 0)}`, ` artifacts:${completed.artifacts.length ? "\n " + completed.artifacts.map((a) => `${a.name} → ${a.uri}`).join("\n ") : " none"}`, p.error ? ` error: ${p.error}` : "", ].filter(Boolean).join("\n"), { job: completed }, ); } return toolText(`Unknown action "${p.action}". Use run | progress | watch | watch_stop | watch_status | finish.`); }, }); // ----------------------------------------------------------------- ledger pi.registerTool({ name: "lvrged_factory_ledger", label: "LVRGED Factory Ledger", description: "Query the job ledger: list jobs with filters (status, deployment, model, workflow, since), and/or roll up cost with group_by (provider | gpu | model | workflow | deployment). Use for 'how much did the last 100 videos cost', 'which jobs used workflow v3', 'which GPU is cheapest for this workload', 'what is queued'.", parameters: Type.Object({ status: Type.Optional(Type.String({ description: "Filter: running | queued | completed | failed" })), deployment_id: Type.Optional(Type.String()), model: Type.Optional(Type.String()), workflow: Type.Optional(Type.String()), workflow_version: Type.Optional(Type.Number()), since: Type.Optional(Type.String({ description: "ISO date or 'today' | 'month'" })), limit: Type.Optional(Type.Number({ description: "Max rows (default 25)", default: 25 })), sum_cost: Type.Optional(Type.Boolean({ description: "Append total cost of the matched set" })), group_by: Type.Optional(Type.String({ description: "Cost rollup instead of a row list: provider | gpu | model | workflow | deployment", enum: ["provider", "gpu", "model", "workflow", "deployment"] })), }), async execute(_id, p: any) { let jobs = loadJobs(); if (p.status) jobs = jobs.filter((j) => j.status === p.status); if (p.deployment_id) jobs = jobs.filter((j) => j.deployment_id === p.deployment_id); if (p.model) jobs = jobs.filter((j) => j.model === p.model); if (p.workflow) jobs = jobs.filter((j) => j.workflow === p.workflow); if (p.workflow_version !== undefined) jobs = jobs.filter((j) => j.workflow_version === p.workflow_version); if (p.since === "today") jobs = jobs.filter((j) => j.ts.slice(0, 10) === new Date().toISOString().slice(0, 10)); else if (p.since === "month") jobs = jobs.filter((j) => j.ts.slice(0, 7) === new Date().toISOString().slice(0, 7)); else if (p.since) jobs = jobs.filter((j) => j.ts >= p.since); const total = jobs.reduce((s, j) => s + (j.cost_usd || 0), 0); if (p.group_by) { const done = jobs.filter((j) => j.status === "completed"); const doneTotal = done.reduce((s, j) => s + (j.cost_usd || 0), 0); const out = [`SPEND${p.since ? ` ${String(p.since).toUpperCase()}` : ""}: ${fmtUsd(doneTotal)} across ${done.length} completed jobs`]; const key = p.group_by === "deployment" ? "deployment_id" : p.group_by; const groups = new Map(); for (const j of done) { const k = String((j as any)[key] || "unknown"); const g = groups.get(k) || { n: 0, cost: 0 }; g.n += 1; g.cost += j.cost_usd || 0; groups.set(k, g); } for (const [k, g] of [...groups.entries()].sort((a, b) => b[1].cost - a[1].cost)) { out.push(` ${k.padEnd(20)} ${g.n.toString().padStart(4)} jobs ${fmtUsd(g.cost)}`); } return toolText(out.join("\n"), { total_cost_usd: Math.round(doneTotal * 1000) / 1000, count: done.length }); } const rows = jobs.slice(-(p.limit || 25)).reverse(); const out = rows.map((j) => `${j.id} ${j.status.padEnd(9)} ${(j.workflow || "").slice(0, 22).padEnd(22)} ${j.model || ""} ${j.duration_s ? j.duration_s + "s" : ""} ${j.cost_usd !== undefined ? fmtUsd(j.cost_usd) : ""}${j.artifacts.length ? " → " + j.artifacts[0].uri : ""}` ); if (!rows.length) return toolText(`No jobs match.`); const lines = [...out]; if (p.sum_cost) lines.push(`TOTAL (${jobs.length} jobs): ${fmtUsd(total)}`); return toolText(lines.join("\n"), { count: jobs.length, total_cost_usd: Math.round(total * 1000) / 1000 }); }, }); // --------------------------------------------------------------- manifest pi.registerTool({ name: "lvrged_factory_manifest", label: "LVRGED Factory Manifest", description: "Add or update a registry manifest. kind=model → models.json (VRAM requirements, source, runtime; use when discovering a model). kind=workflow → workflows.json (executable ComfyUI/HTTP/provider-API workflow so lvrged_factory_job action=run can target it by id).", parameters: Type.Object({ kind: Type.String({ description: "model | workflow", enum: ["model", "workflow"] }), id: Type.String({ description: "Stable id, e.g. wan-2.2-5b or h3-text-to-video" }), // model name: Type.Optional(Type.String({ description: "model: display name" })), source_url: Type.Optional(Type.String({ description: "model: where the weights live (HF/GitHub)" })), vram_gb_min: Type.Optional(Type.Number({ description: "model: minimum VRAM in GB" })), vram_gb_fp8: Type.Optional(Type.Number({ description: "model: VRAM in GB with fp8/INT8 quantization, if known" })), gpu_class: Type.Optional(Type.String({ description: "model: recommended GPU class, e.g. RTX PRO 6000 96GB" })), runtime: Type.Optional(Type.String({ description: "comfyui (default) | http | provider-api | native" })), notes: Type.Optional(Type.String()), verified: Type.Optional(Type.Boolean({ description: "True once a test generation ran successfully" })), // workflow description: Type.Optional(Type.String({ description: "workflow: what it does" })), model: Type.Optional(Type.String({ description: "workflow: model id it belongs to" })), input_schema: Type.Optional(Type.Object({}, { additionalProperties: true, description: "workflow: input field name → type" })), }), async execute(_id, p: any) { const fs = await import("node:fs"); if (p.kind === "model") { if (!p.name || !p.source_url || p.vram_gb_min === undefined || !p.gpu_class) { return toolText("kind=model needs name, source_url, vram_gb_min, gpu_class."); } const models = loadModels(); const man: ModelManifest = { id: p.id, name: p.name, source_url: p.source_url, vram_gb_min: p.vram_gb_min, vram_gb_fp8: p.vram_gb_fp8, gpu_class: p.gpu_class, runtime: p.runtime || "comfyui", notes: p.notes || "", verified: p.verified || false }; const i = models.findIndex((m) => m.id === p.id); if (i >= 0) models[i] = man; else models.push(man); fs.writeFileSync(`${stateDir()}/models.json`, JSON.stringify(models, null, 2) + "\n"); return toolText(`Model ${p.id} saved (${p.name}, ≥${p.vram_gb_min}GB VRAM).`); } if (p.kind === "workflow") { if (!p.runtime || !p.description) return toolText("kind=workflow needs runtime and description."); const workflows = loadWorkflows(); const w: WorkflowManifest = { id: p.id, runtime: p.runtime, description: p.description, model: p.model, input_schema: p.input_schema || {}, verified: false }; const i = workflows.findIndex((x) => x.id === p.id); if (i >= 0) workflows[i] = w; else workflows.push(w); fs.writeFileSync(`${stateDir()}/workflows.json`, JSON.stringify(workflows, null, 2) + "\n"); return toolText(`Workflow ${p.id} saved (runtime: ${p.runtime}).`); } return toolText(`Unknown kind "${p.kind}". Use model | workflow.`); }, }); // ----------------------------------------------------------------- policy pi.registerTool({ name: "lvrged_factory_policy", label: "LVRGED Factory Policy", description: "Read or update the spend policy: per-job ceiling, daily ceiling, monthly ceiling, and the threshold above which provisioning/destruction requires user confirmation. The agent cannot silently override these — they are the trust boundary for real-money actions (which is why this stays its own tool).", parameters: Type.Object({ ceiling_per_job_usd: Type.Optional(Type.Number()), ceiling_daily_usd: Type.Optional(Type.Number()), ceiling_monthly_usd: Type.Optional(Type.Number()), confirm_above_usd: Type.Optional(Type.Number()), idle_shutdown_after_min: Type.Optional(Type.Number({ description: "Minutes of idle before the agent should shut the machine down" })), }), async execute(_id, p: any) { const policy = loadPolicy(); const updated = { ...policy, ...Object.fromEntries(Object.entries(p).filter(([, v]) => v !== undefined)) }; savePolicy(updated); return toolText( [ "SPEND POLICY", ` per-job ceiling: ${fmtUsd(updated.ceiling_per_job_usd)}`, ` daily ceiling: ${fmtUsd(updated.ceiling_daily_usd)}`, ` monthly ceiling: ${fmtUsd(updated.ceiling_monthly_usd)}`, ` confirm above: ${fmtUsd(updated.confirm_above_usd)}`, ` idle shutdown: after ${updated.idle_shutdown_after_min} min`, ].join("\n"), { policy: updated }, ); }, }); } // --------------------------------------------------------------------------- // Helpers // --------------------------------------------------------------------------- function estimateHourly(provider: string, gpu: string, cloud: string = "SECURE"): number { // Reference rates from docs/pricing snapshots + session-verified 2026-08-12 // (`runpodctl gpu list`: PRO 6000 community 1.69 / secure 2.09). The agent // verifies live rates at deploy time — these only feed confirmation gates. const secure: Record> = { runpod: { "RTX PRO 6000": 2.09, "RTX 5090": 0.99, "RTX 4090": 0.69, "RTX 3090": 0.46, "L40S": 0.99, "RTX A6000": 0.49, "H100 SXM": 2.99, "H100 PCIe": 2.89, "A100 SXM": 1.49, "H200": 4.39, "B200": 5.89, "RTX PRO 4500": 0.72, "L4": 0.44 }, }; const community: Record> = { runpod: { "RTX PRO 6000": 1.69, "RTX 5090": 0.69, "RTX 4090": 0.34, "RTX 3090": 0.22, "L40S": 0.79, "RTX A6000": 0.33, "L4": 0.44 }, }; const table = cloud === "COMMUNITY" ? community : secure; const g = gpu.toUpperCase(); const row = table[provider] || secure[provider] || {}; const exact = Object.entries(row).find(([k]) => g.includes(k.toUpperCase())); if (exact) return exact[1]; return g.includes("5090") || g.includes("H100") || g.includes("B200") || g.includes("H200") ? 1.5 : 0.5; } function defaultGpuForVram(_vram: number): string { return "RTX PRO 6000"; // the one lane: 96GB covers every model in the registry, fp16 H3 included } function cheapestProviderFor(_gpu: string): string { return "runpod"; // the one first-class adapter }