| 1 | /** |
| 2 | * An agent's own computer (docs.g1t.sh/guides/agents/, "Its computer"): |
| 3 | * one Durable Object per workspace agent, named `computer:<agent_id>`, in |
| 4 | * front of one container whose process is the runner in `MODE=computer` |
| 5 | * (crates/runner computer.rs). |
| 6 | * |
| 7 | * It is a persistent home, not an always-on machine. Awake, the container |
| 8 | * runs and its seconds are metered as agent sandbox time, attributed to |
| 9 | * the agent and to whoever asked for the session that woke it. Ten minutes |
| 10 | * after the last command with nothing running, the home (`/home/agent`) is |
| 11 | * saved to object storage as one tar+zstd (`HOMES`, R2 on g1t's cloud, |
| 12 | * streamed up in parts, never buffered whole) and the container stops. |
| 13 | * Waking restores the home into a fresh container. Reset wipes the home; |
| 14 | * an archived agent's home is kept `COMPUTER_KEEP_DAYS` and then deleted |
| 15 | * by an alarm. |
| 16 | * |
| 17 | * Without `HOMES` (the bucket not made yet on an installation) everything |
| 18 | * works, but the home is forgotten at every sleep, and the status says so |
| 19 | * (`disk_attached: false`). |
| 20 | * |
| 21 | * Every wake, sleep and reset is written to the workspace's audit log, as |
| 22 | * the agents service writes its own entries. |
| 23 | */ |
| 24 | import { Container, type StopParams } from "@cloudflare/containers"; |
| 25 | |
| 26 | import { |
| 27 | type AgentComputerCommand, |
| 28 | type AgentComputerState, |
| 29 | type AgentComputerStatus, |
| 30 | type ComputerExecArgs, |
| 31 | type ComputerExecResult, |
| 32 | type ComputerWakeArgs, |
| 33 | type Result, |
| 34 | COMPUTER_DISK_CAP_BYTES, |
| 35 | COMPUTER_IDLE_MINUTES, |
| 36 | COMPUTER_KEEP_DAYS, |
| 37 | COMPUTER_TRANSCRIPT_BYTES, |
| 38 | COMPUTER_TRANSCRIPTS_KEPT, |
| 39 | ComputeGate, |
| 40 | actualMicros, |
| 41 | billingClient, |
| 42 | fail, |
| 43 | ok, |
| 44 | sandboxEstimateMicros, |
| 45 | } from "@g1t/contracts"; |
| 46 | |
| 47 | import { concat, homeKey, ndjson, transcriptOf } from "./computer-shape"; |
| 48 | import { describeError } from "./lifecycle"; |
| 49 | import type { RunnerEnv } from "./index"; |
| 50 | |
| 51 | export { gather, homeKey, outcomeLine } from "./computer-shape"; |
| 52 | |
| 53 | /** The port the container's supervisor listens on (crates/runner computer.rs `PORT`). */ |
| 54 | export const COMPUTER_PORT = 8787; |
| 55 | /** How long a wake is reserved for with billing before it starts: the first stretch awake. */ |
| 56 | const RESERVE_MINUTES = 30; |
| 57 | /** Snapshots over the cap are refused; this many refusals in a row and the computer stops unsaved. */ |
| 58 | const MAX_FAILED_SNAPSHOTS = 6; |
| 59 | /** Each part of a snapshot's multipart upload (R2 asks for at least 5 MiB, the last part excepted). */ |
| 60 | const PART_BYTES = 8 * 1024 * 1024; |
| 61 | /** How long one command may be waited on end to end, past its own timeout. */ |
| 62 | const EXEC_GRACE_MS = 15_000; |
| 63 | |
| 64 | /** What billing reserved for a stretch awake, settled when it sleeps. */ |
| 65 | type Held = { id: string; workspace: string; microsPerSecond: number }; |
| 66 | |
| 67 | /** Everything the object keeps about its computer. */ |
| 68 | type Kept = { |
| 69 | agentId: string; |
| 70 | workspace: string | null; |
| 71 | agentHandle: string | null; |
| 72 | askedBy: string | null; |
| 73 | state: AgentComputerState; |
| 74 | since: string; |
| 75 | lastWokeAt: string | null; |
| 76 | lastSleptAt: string | null; |
| 77 | snapshotBytes: number | null; |
| 78 | snapshotAt: string | null; |
| 79 | diskUsedBytes: number; |
| 80 | problem: string | null; |
| 81 | deleteAfter: string | null; |
| 82 | /** When the current stretch awake started, for the meter; null when not awake. */ |
| 83 | meterStarted: number | null; |
| 84 | reservation: Held | null; |
| 85 | token: string | null; |
| 86 | failedSnapshots: number; |
| 87 | }; |
| 88 | |
| 89 | const iso = () => new Date().toISOString(); |
| 90 | |
| 91 | /** One gate per isolate, as the sandboxes keep. */ |
| 92 | let gate: ComputeGate | null = null; |
| 93 | |
| 94 | /** The computer's status as reported, from what is kept. */ |
| 95 | export function statusOf(kept: Kept, attached: boolean): AgentComputerStatus { |
| 96 | return { |
| 97 | agent_id: kept.agentId, |
| 98 | state: kept.state, |
| 99 | since: kept.since, |
| 100 | disk_used_bytes: kept.diskUsedBytes, |
| 101 | disk_cap_bytes: COMPUTER_DISK_CAP_BYTES, |
| 102 | last_woke_at: kept.lastWokeAt, |
| 103 | last_slept_at: kept.lastSleptAt, |
| 104 | snapshot_bytes: kept.snapshotBytes, |
| 105 | snapshot_at: kept.snapshotAt, |
| 106 | where: "g1t_cloud", |
| 107 | disk_attached: attached, |
| 108 | problem: kept.problem, |
| 109 | delete_after: kept.deleteAfter, |
| 110 | }; |
| 111 | } |
| 112 | |
| 113 | export class AgentComputer extends Container<RunnerEnv> { |
| 114 | defaultPort = COMPUTER_PORT; |
| 115 | // Idle: this long after the last request to the container (every command |
| 116 | // is one), `onActivityExpired` saves the home and stops it. |
| 117 | sleepAfter = `${COMPUTER_IDLE_MINUTES}m`; |
| 118 | |
| 119 | /** Commands running now: a busy computer is not put to sleep for idleness. */ |
| 120 | private running = 0; |
| 121 | /** A wake in progress, so two commands arriving together start one container. */ |
| 122 | private waking: Promise<Result<AgentComputerStatus>> | null = null; |
| 123 | |
| 124 | private get attached(): boolean { |
| 125 | return Boolean(this.env.HOMES); |
| 126 | } |
| 127 | |
| 128 | private async kept(agentId: string): Promise<Kept> { |
| 129 | const found = await this.ctx.storage.get<Kept>("computer"); |
| 130 | if (found) return found; |
| 131 | const fresh: Kept = { |
| 132 | agentId, |
| 133 | workspace: null, |
| 134 | agentHandle: null, |
| 135 | askedBy: null, |
| 136 | state: "asleep", |
| 137 | since: iso(), |
| 138 | lastWokeAt: null, |
| 139 | lastSleptAt: null, |
| 140 | snapshotBytes: null, |
| 141 | snapshotAt: null, |
| 142 | diskUsedBytes: 0, |
| 143 | problem: null, |
| 144 | deleteAfter: null, |
| 145 | meterStarted: null, |
| 146 | reservation: null, |
| 147 | token: null, |
| 148 | failedSnapshots: 0, |
| 149 | }; |
| 150 | await this.ctx.storage.put("computer", fresh); |
| 151 | return fresh; |
| 152 | } |
| 153 | |
| 154 | private async save(kept: Kept): Promise<void> { |
| 155 | await this.ctx.storage.put("computer", kept); |
| 156 | } |
| 157 | |
| 158 | /** A request to the container's supervisor, with its token. */ |
| 159 | private async call(kept: Kept, path: string, init: RequestInit = {}): Promise<Response> { |
| 160 | const headers = new Headers(init.headers ?? {}); |
| 161 | if (kept.token) headers.set("authorization", `Bearer ${kept.token}`); |
| 162 | return this.containerFetch(`http://computer${path}`, { ...init, headers }, COMPUTER_PORT); |
| 163 | } |
| 164 | |
| 165 | // ── Status ────────────────────────────────────────────────────────────── |
| 166 | |
| 167 | async status(agentId: string): Promise<AgentComputerStatus> { |
| 168 | const kept = await this.kept(agentId); |
| 169 | if (kept.state === "awake") { |
| 170 | const health = await this.call(kept, "/health").then((r) => (r.ok ? r.json<{ disk_used_bytes?: number }>() : null)).catch(() => null); |
| 171 | if (health && typeof health.disk_used_bytes === "number") { |
| 172 | kept.diskUsedBytes = health.disk_used_bytes; |
| 173 | await this.save(kept); |
| 174 | } |
| 175 | } |
| 176 | return statusOf(kept, this.attached); |
| 177 | } |
| 178 | |
| 179 | async commands(agentId: string, sessionId: string | null = null): Promise<AgentComputerCommand[]> { |
| 180 | await this.kept(agentId); |
| 181 | const all = (await this.ctx.storage.get<AgentComputerCommand[]>("commands")) ?? []; |
| 182 | return sessionId ? all.filter((command) => command.session_id === sessionId) : all; |
| 183 | } |
| 184 | |
| 185 | // ── Waking ────────────────────────────────────────────────────────────── |
| 186 | |
| 187 | async wake(args: ComputerWakeArgs): Promise<Result<AgentComputerStatus>> { |
| 188 | if (this.waking) return this.waking; |
| 189 | this.waking = this.doWake(args).finally(() => { |
| 190 | this.waking = null; |
| 191 | }); |
| 192 | return this.waking; |
| 193 | } |
| 194 | |
| 195 | private async doWake(args: ComputerWakeArgs): Promise<Result<AgentComputerStatus>> { |
| 196 | const kept = await this.kept(args.agent_id); |
| 197 | if (kept.state === "awake") return ok(statusOf(kept, this.attached)); |
| 198 | if (kept.state === "sleeping") return fail("conflict", "The computer is going to sleep; try again in a moment."); |
| 199 | const workspace = args.workspace.toLowerCase(); |
| 200 | // Its workspace's plan says whether it may run, and what is reserved. |
| 201 | gate ??= new ComputeGate(this.env.BILLING); |
| 202 | const [microsPerSecond, ent] = await Promise.all([gate.microsPerSecond(), gate.entitlements(workspace)]); |
| 203 | const admission = await gate.admit( |
| 204 | { |
| 205 | workspace, |
| 206 | repo: { namespace: workspace, name: `agents/${args.agent_handle}` }, |
| 207 | public: false, |
| 208 | kind: "agent", |
| 209 | estimateMicros: sandboxEstimateMicros(RESERVE_MINUTES, microsPerSecond), |
| 210 | // Its model is metered by the agents service; this is machine time only. |
| 211 | hostedModel: false, |
| 212 | }, |
| 213 | ent, |
| 214 | ); |
| 215 | if (!admission.ok) return fail("payment_required", admission.message); |
| 216 | kept.workspace = workspace; |
| 217 | kept.agentHandle = args.agent_handle; |
| 218 | kept.askedBy = args.asked_by ?? null; |
| 219 | kept.reservation = admission.reservation ? { id: admission.reservation.id, workspace, microsPerSecond } : null; |
| 220 | kept.token = crypto.randomUUID(); |
| 221 | kept.state = "waking"; |
| 222 | kept.since = iso(); |
| 223 | kept.problem = null; |
| 224 | await this.save(kept); |
| 225 | try { |
| 226 | await this.startAndWaitForPorts(COMPUTER_PORT, undefined, { |
| 227 | envVars: { MODE: "computer", G1T_COMPUTER_TOKEN: kept.token, ...(this.env.ABUSE_WATCH === "off" ? { G1T_ABUSE: "off" } : {}) }, |
| 228 | enableInternet: true, |
| 229 | }); |
| 230 | } catch (error) { |
| 231 | console.error("computer not started", args.agent_id, describeError(error)); |
| 232 | await this.release(kept); |
| 233 | kept.state = "asleep"; |
| 234 | kept.since = iso(); |
| 235 | kept.problem = `The computer could not start: ${describeError(error)}`; |
| 236 | await this.save(kept); |
| 237 | return fail("unavailable", "The computer could not start just now. Try again in a moment."); |
| 238 | } |
| 239 | // The meter runs from the moment the container is up, restore included. |
| 240 | kept.meterStarted = Date.now(); |
| 241 | await this.save(kept); |
| 242 | await this.restore(kept); |
| 243 | kept.state = "awake"; |
| 244 | kept.since = iso(); |
| 245 | kept.lastWokeAt = kept.since; |
| 246 | await this.save(kept); |
| 247 | this.audit(kept, "computer_wake", `Woke @${kept.agentHandle}'s computer${kept.snapshotAt ? ", restoring its home" : " with a fresh home"}.`); |
| 248 | return ok(await this.status(args.agent_id)); |
| 249 | } |
| 250 | |
| 251 | /** Puts the saved home back into a container that just started. */ |
| 252 | private async restore(kept: Kept): Promise<void> { |
| 253 | if (!this.env.HOMES || !kept.snapshotAt) return; |
| 254 | try { |
| 255 | const object = await this.env.HOMES.get(homeKey(kept.agentId)); |
| 256 | if (!object) { |
| 257 | kept.snapshotAt = null; |
| 258 | kept.snapshotBytes = null; |
| 259 | return; |
| 260 | } |
| 261 | const response = await this.call(kept, "/restore", { method: "POST", body: object.body, headers: { "content-type": "application/zstd" } }); |
| 262 | if (!response.ok) throw new Error(`restore answered ${response.status}: ${(await response.text().catch(() => "")).slice(0, 200)}`); |
| 263 | const answer = await response.json<{ disk_used_bytes?: number }>().catch((): { disk_used_bytes?: number } => ({})); |
| 264 | if (typeof answer.disk_used_bytes === "number") kept.diskUsedBytes = answer.disk_used_bytes; |
| 265 | } catch (error) { |
| 266 | console.error("computer home not restored", kept.agentId, describeError(error)); |
| 267 | kept.problem = "Its saved home could not be restored this time, so it woke with a fresh one. The saved home is kept for the next wake."; |
| 268 | } |
| 269 | } |
| 270 | |
| 271 | /** A wake that admitted and then failed gives back what was reserved. */ |
| 272 | private async release(kept: Kept): Promise<void> { |
| 273 | if (!kept.reservation) return; |
| 274 | gate ??= new ComputeGate(this.env.BILLING); |
| 275 | await gate.settle(kept.reservation.id, 0); |
| 276 | kept.reservation = null; |
| 277 | } |
| 278 | |
| 279 | /** Awake, or why not. */ |
| 280 | private async awake(args: ComputerWakeArgs): Promise<Result<Kept>> { |
| 281 | const kept = await this.kept(args.agent_id); |
| 282 | if (kept.state !== "awake") { |
| 283 | const woke = await this.wake(args); |
| 284 | if (!woke.ok) return woke; |
| 285 | } |
| 286 | return ok(await this.kept(args.agent_id)); |
| 287 | } |
| 288 | |
| 289 | // ── Commands and files ───────────────────────────────────────────────── |
| 290 | |
| 291 | async exec(args: ComputerExecArgs): Promise<Result<ComputerExecResult>> { |
| 292 | const cmd = String(args.cmd ?? "").trim(); |
| 293 | if (!cmd) return fail("invalid", "Give the command to run."); |
| 294 | const ready = await this.awake(args); |
| 295 | if (!ready.ok) return ready; |
| 296 | const kept = ready.value; |
| 297 | const timeout = Math.max(1, Math.min(1800, Math.floor(Number(args.timeout_seconds ?? 120)) || 120)); |
| 298 | const startedAt = iso(); |
| 299 | this.running++; |
| 300 | const lines: { stream: string; line: string }[] = []; |
| 301 | let closing: { exit_code?: number; duration_ms?: number; truncated?: boolean; timed_out?: boolean; note?: string } | null = null; |
| 302 | let problem: string | null = null; |
| 303 | try { |
| 304 | const controller = new AbortController(); |
| 305 | const timer = setTimeout(() => controller.abort(), timeout * 1000 + EXEC_GRACE_MS); |
| 306 | try { |
| 307 | const response = await this.call(kept, "/exec", { |
| 308 | method: "POST", |
| 309 | headers: { "content-type": "application/json" }, |
| 310 | body: JSON.stringify({ cmd, cwd: args.cwd ?? null, timeout_seconds: timeout }), |
| 311 | signal: controller.signal, |
| 312 | }); |
| 313 | if (!response.ok || !response.body) { |
| 314 | const text = await response.text().catch(() => ""); |
| 315 | let message = text; |
| 316 | try { |
| 317 | message = (JSON.parse(text) as { message?: string }).message ?? text; |
| 318 | } catch { |
| 319 | // Plain text, kept as it is. |
| 320 | } |
| 321 | return fail("invalid", message || `The computer answered ${response.status}.`); |
| 322 | } |
| 323 | for await (const line of ndjson(response.body)) { |
| 324 | if (typeof line.exit_code === "number") closing = line as unknown as typeof closing; |
| 325 | else if (typeof line.line === "string") lines.push({ stream: String(line.stream ?? "stdout"), line: line.line }); |
| 326 | } |
| 327 | } finally { |
| 328 | clearTimeout(timer); |
| 329 | } |
| 330 | } catch (error) { |
| 331 | problem = describeError(error); |
| 332 | console.error("computer exec failed", kept.agentId, problem); |
| 333 | } finally { |
| 334 | this.running--; |
| 335 | } |
| 336 | const shaped = transcriptOf(lines, closing, problem, COMPUTER_TRANSCRIPT_BYTES); |
| 337 | const command: AgentComputerCommand = { |
| 338 | id: `cmd_${crypto.randomUUID().replace(/-/g, "").slice(0, 20)}`, |
| 339 | session_id: args.session_id ?? null, |
| 340 | asked_by: args.asked_by ?? kept.askedBy, |
| 341 | started_at: startedAt, |
| 342 | cmd, |
| 343 | cwd: args.cwd?.trim() || "/home/agent", |
| 344 | exit_code: shaped.exit_code, |
| 345 | duration_ms: shaped.duration_ms ?? Math.max(0, Date.now() - Date.parse(startedAt)), |
| 346 | output: shaped.output, |
| 347 | truncated: shaped.truncated, |
| 348 | timed_out: shaped.timed_out, |
| 349 | }; |
| 350 | const all = (await this.ctx.storage.get<AgentComputerCommand[]>("commands")) ?? []; |
| 351 | await this.ctx.storage.put("commands", [command, ...all].slice(0, COMPUTER_TRANSCRIPTS_KEPT)); |
| 352 | return ok({ command, status: await this.status(kept.agentId) }); |
| 353 | } |
| 354 | |
| 355 | async readFile(args: ComputerWakeArgs & { path: string }): Promise<Result<{ path: string; text: string; bytes: number }>> { |
| 356 | const path = String(args.path ?? "").trim(); |
| 357 | if (!path) return fail("invalid", "Give the file's path under the home."); |
| 358 | const ready = await this.awake(args); |
| 359 | if (!ready.ok) return ready; |
| 360 | const response = await this.call(ready.value, `/files?path=${encodeURIComponent(path)}`).catch((error: unknown) => { |
| 361 | console.error("computer read failed", args.agent_id, describeError(error)); |
| 362 | return null; |
| 363 | }); |
| 364 | if (!response) return fail("unavailable", "The computer didn't answer."); |
| 365 | if (response.status === 404) return fail("not_found", `There is no file at ${path}.`); |
| 366 | if (!response.ok) return fail("invalid", await said(response)); |
| 367 | const bytes = new Uint8Array(await response.arrayBuffer()); |
| 368 | return ok({ path, text: new TextDecoder().decode(bytes), bytes: bytes.byteLength }); |
| 369 | } |
| 370 | |
| 371 | async writeFile(args: ComputerWakeArgs & { path: string; text: string }): Promise<Result<{ path: string; bytes: number }>> { |
| 372 | const path = String(args.path ?? "").trim(); |
| 373 | if (!path) return fail("invalid", "Give the file's path under the home."); |
| 374 | const ready = await this.awake(args); |
| 375 | if (!ready.ok) return ready; |
| 376 | const body = new TextEncoder().encode(String(args.text ?? "")); |
| 377 | const response = await this.call(ready.value, `/files?path=${encodeURIComponent(path)}`, { method: "PUT", body, headers: { "content-type": "application/octet-stream" } }).catch((error: unknown) => { |
| 378 | console.error("computer write failed", args.agent_id, describeError(error)); |
| 379 | return null; |
| 380 | }); |
| 381 | if (!response) return fail("unavailable", "The computer didn't answer."); |
| 382 | if (!response.ok) return fail("invalid", await said(response)); |
| 383 | return ok({ path, bytes: body.byteLength }); |
| 384 | } |
| 385 | |
| 386 | // ── Sleeping ──────────────────────────────────────────────────────────── |
| 387 | |
| 388 | async sleep(agentId: string): Promise<Result<AgentComputerStatus>> { |
| 389 | const kept = await this.kept(agentId); |
| 390 | if (kept.state === "asleep") return ok(statusOf(kept, this.attached)); |
| 391 | if (kept.state !== "awake") return fail("conflict", `The computer is ${kept.state}; try again in a moment.`); |
| 392 | if (this.running > 0) return fail("conflict", "A command is still running. It sleeps once the command ends and it has been idle for ten minutes."); |
| 393 | kept.state = "sleeping"; |
| 394 | kept.since = iso(); |
| 395 | await this.save(kept); |
| 396 | const saved = await this.snapshot(kept); |
| 397 | if (!saved.ok) { |
| 398 | kept.failedSnapshots++; |
| 399 | kept.problem = saved.error.message; |
| 400 | if (kept.failedSnapshots < MAX_FAILED_SNAPSHOTS) { |
| 401 | kept.state = "awake"; |
| 402 | kept.since = iso(); |
| 403 | await this.save(kept); |
| 404 | // Another try at the next idle expiry. |
| 405 | this.renewActivityTimeout(); |
| 406 | return fail("conflict", saved.error.message); |
| 407 | } |
| 408 | kept.problem = `${saved.error.message} It was stopped without saving after ${MAX_FAILED_SNAPSHOTS} tries, so changes since its last save are gone.`; |
| 409 | } else { |
| 410 | kept.failedSnapshots = 0; |
| 411 | kept.problem = null; |
| 412 | } |
| 413 | await this.save(kept); |
| 414 | await this.call(kept, "/stop", { method: "POST" }).catch(() => null); |
| 415 | await this.stop().catch(() => undefined); |
| 416 | await this.stopped(0, "asked"); |
| 417 | this.audit(kept, "computer_sleep", `Put @${kept.agentHandle}'s computer to sleep${saved.ok ? ", home saved" : " without saving its home"}.`); |
| 418 | return ok(statusOf((await this.kept(agentId))!, this.attached)); |
| 419 | } |
| 420 | |
| 421 | /** Streams the home out of the container and up to object storage, in parts. */ |
| 422 | private async snapshot(kept: Kept): Promise<Result<null>> { |
| 423 | const bucket = this.env.HOMES; |
| 424 | if (!bucket) { |
| 425 | // Nothing to save to: the home is forgotten, and the status says so. |
| 426 | kept.snapshotAt = null; |
| 427 | kept.snapshotBytes = null; |
| 428 | return ok(null); |
| 429 | } |
| 430 | let response: Response; |
| 431 | try { |
| 432 | response = await this.call(kept, "/snapshot", { method: "POST" }); |
| 433 | } catch (error) { |
| 434 | return fail("unavailable", `The home could not be read for saving: ${describeError(error)}`); |
| 435 | } |
| 436 | if (response.status === 413) return fail("invalid", await said(response)); |
| 437 | if (!response.ok || !response.body) return fail("unavailable", `The home could not be saved: the computer answered ${response.status}.`); |
| 438 | const upload = await bucket.createMultipartUpload(homeKey(kept.agentId)); |
| 439 | const parts: R2UploadedPart[] = []; |
| 440 | let total = 0; |
| 441 | try { |
| 442 | let held: Uint8Array[] = []; |
| 443 | let size = 0; |
| 444 | const flush = async () => { |
| 445 | const part = concat(held, size); |
| 446 | held = []; |
| 447 | size = 0; |
| 448 | parts.push(await upload.uploadPart(parts.length + 1, part)); |
| 449 | }; |
| 450 | const reader = response.body.getReader(); |
| 451 | for (;;) { |
| 452 | const { done, value } = await reader.read(); |
| 453 | if (done) break; |
| 454 | held.push(value); |
| 455 | size += value.byteLength; |
| 456 | total += value.byteLength; |
| 457 | if (size >= PART_BYTES) await flush(); |
| 458 | } |
| 459 | if (size > 0 || parts.length === 0) await flush(); |
| 460 | await upload.complete(parts); |
| 461 | } catch (error) { |
| 462 | await upload.abort().catch(() => undefined); |
| 463 | return fail("unavailable", `The home could not be saved: ${describeError(error)}`); |
| 464 | } |
| 465 | kept.snapshotBytes = total; |
| 466 | kept.snapshotAt = iso(); |
| 467 | const used = Number(response.headers.get("x-g1t-disk-used")); |
| 468 | if (Number.isFinite(used) && used > 0) kept.diskUsedBytes = used; |
| 469 | return ok(null); |
| 470 | } |
| 471 | |
| 472 | /** |
| 473 | * Everything the end of a stretch awake does, once: the meter, the |
| 474 | * reservation, the state. From `sleep`, and from `onStop` when the |
| 475 | * container went by itself. |
| 476 | */ |
| 477 | private async stopped(exitCode: number, how: "asked" | "exit"): Promise<void> { |
| 478 | const kept = await this.ctx.storage.get<Kept>("computer"); |
| 479 | if (!kept) return; |
| 480 | if (kept.meterStarted != null) { |
| 481 | const seconds = Math.max(1, Math.ceil((Date.now() - kept.meterStarted) / 1000)); |
| 482 | const recorded = await billingClient(this.env.BILLING) |
| 483 | .recordSandbox({ |
| 484 | workspace: kept.workspace ?? "", |
| 485 | seconds, |
| 486 | description: `@${kept.agentHandle ?? "agent"}'s computer, awake`, |
| 487 | repo: null, |
| 488 | reference: `computer/${kept.agentId}/${kept.meterStarted}`, |
| 489 | kind: "agent", |
| 490 | selfHosted: false, |
| 491 | instance: null, |
| 492 | agent: kept.agentHandle, |
| 493 | askedBy: kept.askedBy, |
| 494 | }) |
| 495 | .catch((error: unknown) => ({ ok: false as const, error: { message: String(error) } })); |
| 496 | if (!recorded.ok) console.log("computer time not recorded", kept.workspace, seconds, recorded.error.message); |
| 497 | if (kept.reservation) { |
| 498 | gate ??= new ComputeGate(this.env.BILLING); |
| 499 | await gate.settle(kept.reservation.id, actualMicros(seconds, kept.reservation.microsPerSecond, 0)); |
| 500 | } |
| 501 | kept.meterStarted = null; |
| 502 | kept.reservation = null; |
| 503 | } |
| 504 | if (how === "exit" && (kept.state === "awake" || kept.state === "waking")) { |
| 505 | kept.problem = `The computer stopped by itself (exit ${exitCode}); changes since its last save are gone.`; |
| 506 | } |
| 507 | if (kept.state !== "asleep") { |
| 508 | kept.state = "asleep"; |
| 509 | kept.since = iso(); |
| 510 | kept.lastSleptAt = kept.since; |
| 511 | } |
| 512 | kept.token = null; |
| 513 | await this.save(kept); |
| 514 | } |
| 515 | |
| 516 | /** |
| 517 | * Idle for `sleepAfter`: save and stop, unless a command is running, in |
| 518 | * which case look again later. Never `super`, which would stop the |
| 519 | * container without saving. |
| 520 | */ |
| 521 | override async onActivityExpired(): Promise<void> { |
| 522 | const kept = await this.ctx.storage.get<Kept>("computer"); |
| 523 | if (!kept || kept.state !== "awake") return; |
| 524 | if (this.running > 0) { |
| 525 | this.renewActivityTimeout(); |
| 526 | return; |
| 527 | } |
| 528 | const slept = await this.sleep(kept.agentId); |
| 529 | if (!slept.ok) console.log("computer not put to sleep", kept.agentId, slept.error.message); |
| 530 | } |
| 531 | |
| 532 | override async onStop({ exitCode, reason }: StopParams): Promise<void> { |
| 533 | const kept = await this.ctx.storage.get<Kept>("computer"); |
| 534 | if (kept?.state === "asleep" && kept.meterStarted == null) return; |
| 535 | console.log("computer stopped", kept?.agentId, "exit", exitCode, reason); |
| 536 | await this.stopped(exitCode, "exit"); |
| 537 | } |
| 538 | |
| 539 | override onError(error: unknown): void { |
| 540 | console.error("computer container error", this.ctx.id.toString(), describeError(error)); |
| 541 | } |
| 542 | |
| 543 | override async alarm(alarmProps?: AlarmInvocationInfo): Promise<void> { |
| 544 | try { |
| 545 | await super.alarm(alarmProps); |
| 546 | } catch (error) { |
| 547 | console.error("computer alarm failed", this.ctx.id.toString(), `retry ${alarmProps?.retryCount ?? 0}`, describeError(error)); |
| 548 | throw error; |
| 549 | } |
| 550 | } |
| 551 | |
| 552 | // ── Reset and forgetting ─────────────────────────────────────────────── |
| 553 | |
| 554 | async reset(agentId: string): Promise<Result<AgentComputerStatus>> { |
| 555 | const kept = await this.kept(agentId); |
| 556 | if (kept.state !== "asleep") { |
| 557 | await this.destroy().catch((error: unknown) => console.log("computer not destroyed for reset", agentId, describeError(error))); |
| 558 | await this.stopped(0, "asked"); |
| 559 | } |
| 560 | if (this.env.HOMES) await this.env.HOMES.delete(homeKey(agentId)).catch((error: unknown) => console.log("computer home not deleted", agentId, describeError(error))); |
| 561 | await this.ctx.storage.delete("commands"); |
| 562 | const fresh = (await this.kept(agentId))!; |
| 563 | fresh.snapshotAt = null; |
| 564 | fresh.snapshotBytes = null; |
| 565 | fresh.diskUsedBytes = 0; |
| 566 | fresh.problem = null; |
| 567 | fresh.failedSnapshots = 0; |
| 568 | fresh.state = "asleep"; |
| 569 | fresh.since = iso(); |
| 570 | await this.save(fresh); |
| 571 | this.audit(fresh, "computer_reset", `Reset @${fresh.agentHandle ?? "agent"}'s computer: its home was wiped.`); |
| 572 | return ok(statusOf(fresh, this.attached)); |
| 573 | } |
| 574 | |
| 575 | /** The agent was archived: its home is kept `COMPUTER_KEEP_DAYS`, then deleted. */ |
| 576 | async forget(agentId: string): Promise<Result<AgentComputerStatus>> { |
| 577 | const kept = await this.kept(agentId); |
| 578 | if (kept.state === "awake") await this.sleep(agentId); |
| 579 | const fresh = (await this.kept(agentId))!; |
| 580 | fresh.deleteAfter = new Date(Date.now() + COMPUTER_KEEP_DAYS * 24 * 60 * 60 * 1000).toISOString(); |
| 581 | await this.save(fresh); |
| 582 | await this.schedule(COMPUTER_KEEP_DAYS * 24 * 60 * 60, "purge"); |
| 583 | return ok(statusOf(fresh, this.attached)); |
| 584 | } |
| 585 | |
| 586 | /** The alarm after `forget`: the disk and everything kept here go. */ |
| 587 | async purge(): Promise<void> { |
| 588 | const kept = await this.ctx.storage.get<Kept>("computer"); |
| 589 | if (!kept?.deleteAfter || Date.parse(kept.deleteAfter) > Date.now()) return; |
| 590 | if (kept.state !== "asleep") await this.destroy().catch(() => undefined); |
| 591 | if (this.env.HOMES) await this.env.HOMES.delete(homeKey(kept.agentId)).catch(() => undefined); |
| 592 | await this.ctx.storage.deleteAll(); |
| 593 | } |
| 594 | |
| 595 | /** An entry in the workspace's audit log, as the agents service writes them. Never fails the call. */ |
| 596 | private audit(kept: Kept, action: string, message: string): void { |
| 597 | if (!this.env.EVENTS || !kept.workspace) return; |
| 598 | const entry = { |
| 599 | actorKind: "agent", |
| 600 | actor: kept.agentHandle ?? kept.agentId, |
| 601 | actorId: kept.agentId, |
| 602 | agent: kept.agentHandle, |
| 603 | onBehalfOf: kept.askedBy, |
| 604 | runId: null, |
| 605 | runKind: null, |
| 606 | credentialId: null, |
| 607 | action, |
| 608 | surface: "web", |
| 609 | workspace: kept.workspace, |
| 610 | repo: null, |
| 611 | number: null, |
| 612 | gitRef: null, |
| 613 | path: `agents/${kept.agentHandle ?? kept.agentId}/computer`, |
| 614 | outcome: "allowed", |
| 615 | rule: "owner", |
| 616 | result: "ok", |
| 617 | message, |
| 618 | requestId: `req_${crypto.randomUUID()}`, |
| 619 | }; |
| 620 | this.ctx.waitUntil( |
| 621 | this.env.EVENTS.fetch("https://service/rpc/audit_record", { method: "POST", headers: { "content-type": "application/json" }, body: JSON.stringify({ entries: [entry] }) }) |
| 622 | .then((response) => { |
| 623 | if (!response.ok) throw new Error(`status ${response.status}`); |
| 624 | }) |
| 625 | .catch((error: unknown) => console.error("computer audit entry not recorded", action, kept.agentId, String(error))), |
| 626 | ); |
| 627 | } |
| 628 | } |
| 629 | |
| 630 | /** A JSON error body's `message`, or the status. */ |
| 631 | async function said(response: Response): Promise<string> { |
| 632 | const text = await response.text().catch(() => ""); |
| 633 | try { |
| 634 | return (JSON.parse(text) as { message?: string }).message ?? `The computer answered ${response.status}.`; |
| 635 | } catch { |
| 636 | return text || `The computer answered ${response.status}.`; |
| 637 | } |
| 638 | } |
| 639 | |