From 99dc8ddda5c823fbe1c7d7016840c8275e75e652 Mon Sep 17 00:00:00 2001 From: Mohamed Boudra Date: Sat, 18 Jul 2026 21:21:05 +0200 Subject: [PATCH] Free resources from idle agents automatically (#2203) * fix(server): release resources held by idle agents Keep unarchived agents resumable while closing their provider runtimes after two minutes. Active agent schedules keep runtimes resident. * fix(server): preserve resumable agent state * fix(server): resume agents when listing commands * fix(server): preserve collected agent interactions --- docs/agent-lifecycle.md | 25 +- docs/data-model.md | 2 +- .../server/src/server/agent/agent-loading.ts | 12 +- .../src/server/agent/agent-manager.test.ts | 328 ++++++++++++++++++ .../server/src/server/agent/agent-manager.ts | 157 ++++++++- .../src/server/agent/agent-response-loop.ts | 2 + .../src/server/agent/agent-sdk-types.ts | 5 +- .../opencode-agent-commands.e2e.test.ts | 22 ++ .../agent/providers/opencode-agent.test.ts | 132 +++---- .../server/agent/providers/opencode-agent.ts | 98 ++++-- .../providers/opencode-server-manager.test.ts | 66 ++-- .../providers/opencode/server-manager.ts | 52 +-- .../providers/opencode/test-server-manager.ts | 8 +- .../test-utils/test-opencode-harness.ts | 6 +- packages/server/src/server/bootstrap.ts | 37 ++ .../src/server/hub/daemon-executions.ts | 6 +- packages/server/src/server/loop-service.ts | 15 +- .../src/server/schedule/service.test.ts | 41 +++ .../server/src/server/schedule/service.ts | 13 + packages/server/src/server/session.ts | 33 +- 20 files changed, 879 insertions(+), 181 deletions(-) diff --git a/docs/agent-lifecycle.md b/docs/agent-lifecycle.md index 2be3c3b55..f5e922830 100644 --- a/docs/agent-lifecycle.md +++ b/docs/agent-lifecycle.md @@ -10,7 +10,26 @@ initializing → idle → running → idle (or error → closed) └────────┘ (agent completes a turn, awaits next prompt) ``` -Each agent in `AgentManager` carries a `lastStatus` of `initializing`, `idle`, `running`, `error`, or `closed`. State transitions persist to disk and stream to subscribed clients via WebSocket. +Each live agent in `AgentManager` carries a `lastStatus` of `initializing`, `idle`, `running`, or `error`. `closed` is the persisted, resumable state for an agent record that has no live provider runtime. State transitions persist to disk and stream to subscribed clients via WebSocket. + +## Runtime residency + +An unarchived agent may be `closed` without being deleted or archived. Closing releases its provider +processes and subscriptions while retaining its Paseo identity, persistence handle, timeline, +workspace, labels, title, usage, attention, timestamps, and parent relationship. Opening or prompting +the agent runs through `ensureAgentLoaded()`, which resumes the durable provider session under the +same Paseo agent ID. Provider history is not appended again when the canonical timeline is already +primed. + +The daemon collects an eligible idle runtime after two minutes and sweeps every 15 seconds. Only +unarchived, non-internal agents that are exactly `idle`, have no active or pending run, replacement, +or permission, and have not been activated during the idle window are eligible. `running`, +`initializing`, and `error` agents stay resident. Subagents are considered independently; collection +does not cascade or change parentage. + +Active schedules targeting an existing agent protect that agent from collection. Paused, completed, +and new-agent schedules do not. A pane may remain open after collection; its next prompt resumes the +runtime. ### Cancellation @@ -39,6 +58,10 @@ The provider still owns the underlying runtime. Paseo keeps an agent record so t Archive is a **soft delete**: the agent record stays on disk with `archivedAt` set, the runtime is closed, and the agent disappears from active lists. Archive is **global** — it lives on the server and propagates to every connected client. +Archive is distinct from runtime collection. Archive sets `archivedAt`, invokes the provider's native +archive hook, and cascades to managed children. Runtime collection does none of those things; it only +releases the live runtime and writes `lastStatus: closed` on the still-active record. + `create_agent_request` can opt an agent into `autoArchive`. In that mode the daemon archives the agent after the first terminal turn event (`turn_completed`, `turn_failed`, or `turn_canceled`). When the agent owns an isolated workspace, auto-archive archives that workspace too; the managed worktree is removed when its final workspace reference is gone. Archiving runs through `AgentManager.archiveAgent` (`packages/server/src/server/agent/agent-manager.ts`): diff --git a/docs/data-model.md b/docs/data-model.md index 0d34d49a6..829ca612a 100644 --- a/docs/data-model.md +++ b/docs/data-model.md @@ -83,7 +83,7 @@ Each agent is stored as a separate JSON file, grouped by project directory. | `lastUserMessageAt` | `string?` (ISO 8601) | Last user message timestamp | | `title` | `string?` | User-visible title | | `labels` | `Record` | Key-value labels (default `{}`). `paseo.parent-agent-id` is set automatically for agent-scoped creation and removed by detach — see [agent-lifecycle.md](./agent-lifecycle.md) | -| `lastStatus` | `AgentStatus` | One of: `"initializing"`, `"idle"`, `"running"`, `"error"`, `"closed"` | +| `lastStatus` | `AgentStatus` | One of: `"initializing"`, `"idle"`, `"running"`, `"error"`, `"closed"`. `closed` means the record is resumable but has no live provider runtime; archive remains represented separately by `archivedAt`. | | `lastModeId` | `string?` | Last active mode ID | | `config` | `SerializableConfig?` | Agent session configuration (see below) | | `runtimeInfo` | `RuntimeInfo?` | Live runtime state (see below) | diff --git a/packages/server/src/server/agent/agent-loading.ts b/packages/server/src/server/agent/agent-loading.ts index c3a67af2c..01dd9264b 100644 --- a/packages/server/src/server/agent/agent-loading.ts +++ b/packages/server/src/server/agent/agent-loading.ts @@ -20,7 +20,8 @@ export type AgentLoaderManager = Pick< | "getRegisteredProviderIds" | "hydrateTimelineFromProvider" | "resumeAgentFromPersistence" ->; +> & + Partial>; export interface EnsureAgentLoadedDeps { agentManager: AgentLoaderManager; @@ -33,11 +34,18 @@ export async function ensureAgentLoaded( agentId: string, deps: EnsureAgentLoadedDeps, ): Promise { - const existing = deps.agentManager.getAgent(agentId); + await deps.agentManager.waitForAgentClose?.(agentId); + const existing = + deps.agentManager.touchAgentActivity?.(agentId) ?? deps.agentManager.getAgent(agentId); if (existing) { return existing; } + // A close may have started after the first barrier observed no in-flight + // work. Once the live lookup is empty, this second barrier closes that gap + // before storage-backed resume begins. + await deps.agentManager.waitForAgentClose?.(agentId); + const inflight = pendingAgentInitializations.get(agentId); if (inflight) { return inflight; diff --git a/packages/server/src/server/agent/agent-manager.test.ts b/packages/server/src/server/agent/agent-manager.test.ts index 48e32643c..12dea4e47 100644 --- a/packages/server/src/server/agent/agent-manager.test.ts +++ b/packages/server/src/server/agent/agent-manager.test.ts @@ -17,6 +17,7 @@ import { AgentStorage } from "./agent-storage.js"; import { toAgentPayload } from "./agent-projections.js"; import { PARENT_AGENT_ID_LABEL } from "@getpaseo/protocol/agent-labels"; import { formatSystemNotificationPrompt } from "./agent-prompt.js"; +import { ensureAgentLoaded } from "./agent-loading.js"; import type { StoredAgentRecord } from "./agent-storage.js"; import type { AgentClient, @@ -7159,6 +7160,333 @@ test("closeAgent persists one final closed snapshot", async () => { } }); +test("collectIdleAgents releases an idle runtime and resumes the same agent and timeline", async () => { + const workdir = mkdtempSync(join(tmpdir(), "agent-manager-idle-collection-")); + const storage = new AgentStorage(join(workdir, "agents"), logger); + let activeSession: TestAgentSession | null = null; + const client = new (class extends NativeArchiveRecordingClient { + override async createSession(config: AgentSessionConfig): Promise { + activeSession = new TestAgentSession(config); + return activeSession; + } + })(); + const manager = new AgentManager({ + clients: { codex: client }, + registry: storage, + logger, + idFactory: () => "00000000-0000-4000-8000-000000000210", + }); + + try { + const created = await manager.createAgent({ provider: "codex", cwd: workdir }, undefined, { + workspaceId: "workspace-idle-collection", + }); + await manager.appendTimelineItem(created.id, { + type: "user_message", + text: "Keep this timeline", + }); + activeSession?.pushEvent({ + type: "provider_subagent", + provider: "codex", + event: { + type: "upsert", + id: "retained-provider-child", + title: "Retained provider child", + status: "completed", + }, + }); + await manager.flush(); + const timelineBeforeCollection = manager.getTimeline(created.id); + + const collection = await manager.collectIdleAgents({ + cutoff: new Date(Date.now() + 1_000), + protectedAgentIds: new Set(), + }); + + expect(collection).toEqual({ + collected: [ + { + agentId: created.id, + provider: "codex", + sessionId: created.persistence?.sessionId, + }, + ], + failures: [], + }); + expect(manager.getAgent(created.id)).toBeNull(); + expect(client.archivedHandles).toEqual([]); + const stored = await storage.get(created.id); + expect(stored).toMatchObject({ + id: created.id, + lastStatus: "closed", + workspaceId: "workspace-idle-collection", + }); + expect(stored?.archivedAt).toBeFalsy(); + + const resumed = await ensureAgentLoaded(created.id, { + agentManager: manager, + agentStorage: storage, + logger, + }); + + expect(resumed.id).toBe(created.id); + expect(resumed.persistence).toEqual(created.persistence); + expect(manager.getTimeline(created.id)).toEqual(timelineBeforeCollection); + expect(manager.listProviderSubagents(created.id)).toEqual([ + expect.objectContaining({ + id: "retained-provider-child", + title: "Retained provider child", + status: "completed", + }), + ]); + const idleBeforeOpen = resumed.updatedAt; + await ensureAgentLoaded(created.id, { + agentManager: manager, + agentStorage: storage, + logger, + }); + await expect( + manager.collectIdleAgents({ cutoff: idleBeforeOpen, protectedAgentIds: new Set() }), + ).resolves.toMatchObject({ collected: [] }); + await expect(manager.runAgent(created.id, "Continue the same agent")).resolves.toMatchObject({ + finalText: "", + canceled: false, + }); + expect(manager.getAgent(created.id)?.id).toBe(created.id); + } finally { + await manager.flush().catch(() => undefined); + await storage.flush().catch(() => undefined); + rmSync(workdir, { recursive: true, force: true }); + } +}); + +test("archiving an idle-collected parent still cascades to its managed children", async () => { + const workdir = mkdtempSync(join(tmpdir(), "agent-manager-collected-parent-archive-")); + const storage = new AgentStorage(join(workdir, "agents"), logger); + const manager = new AgentManager({ + clients: { codex: new TestAgentClient() }, + registry: storage, + logger, + }); + + try { + const parent = await manager.createAgent( + { provider: "codex", cwd: workdir, title: "Collected parent" }, + undefined, + { workspaceId: undefined }, + ); + const child = await manager.createAgent( + { provider: "codex", cwd: workdir, title: "Managed child" }, + undefined, + { + labels: { [PARENT_AGENT_ID_LABEL]: parent.id }, + workspaceId: undefined, + }, + ); + + await manager.collectIdleAgents({ + cutoff: new Date(Date.now() + 1_000), + protectedAgentIds: new Set([child.id]), + }); + await manager.archiveSnapshot(parent.id, new Date().toISOString()); + + expect((await storage.get(parent.id))?.archivedAt).toEqual(expect.any(String)); + expect((await storage.get(child.id))?.archivedAt).toEqual(expect.any(String)); + expect(manager.getAgent(child.id)).toBeNull(); + } finally { + await manager.flush().catch(() => undefined); + await storage.flush().catch(() => undefined); + rmSync(workdir, { recursive: true, force: true }); + } +}); + +test("collectIdleAgents leaves recent, protected, internal, running, and error agents resident", async () => { + const workdir = mkdtempSync(join(tmpdir(), "agent-manager-idle-eligibility-")); + const client = new (class extends TestAgentClient { + readonly sessions: TestAgentSession[] = []; + + override async createSession(config: AgentSessionConfig): Promise { + const session = new TestAgentSession(config); + this.sessions.push(session); + return session; + } + })(); + const ids = [ + "00000000-0000-4000-8000-000000000211", + "00000000-0000-4000-8000-000000000212", + "00000000-0000-4000-8000-000000000213", + "00000000-0000-4000-8000-000000000214", + "00000000-0000-4000-8000-000000000215", + ]; + const manager = new AgentManager({ + clients: { codex: client }, + logger, + idFactory: () => ids.shift()!, + }); + + try { + const recent = await manager.createAgent({ provider: "codex", cwd: workdir }, undefined, { + workspaceId: undefined, + }); + const protectedAgent = await manager.createAgent( + { provider: "codex", cwd: workdir }, + undefined, + { workspaceId: undefined }, + ); + const internal = await manager.createAgent( + { provider: "codex", cwd: workdir, internal: true }, + undefined, + { workspaceId: undefined }, + ); + const running = await manager.createAgent({ provider: "codex", cwd: workdir }, undefined, { + workspaceId: undefined, + }); + const failed = await manager.createAgent({ provider: "codex", cwd: workdir }, undefined, { + workspaceId: undefined, + }); + client.sessions[3]!.pushEvent({ + type: "turn_started", + provider: "codex", + turnId: "autonomous-running", + }); + client.sessions[4]!.pushEvent({ + type: "turn_failed", + provider: "codex", + turnId: "autonomous-failed", + error: "provider failed", + }); + await manager.flush(); + + const recentSweep = await manager.collectIdleAgents({ + cutoff: new Date(recent.updatedAt.getTime() - 1), + protectedAgentIds: new Set(), + }); + const protectedSweep = await manager.collectIdleAgents({ + cutoff: new Date(Date.now() + 1_000), + protectedAgentIds: new Set([protectedAgent.id, recent.id]), + }); + + expect(recentSweep.collected).toEqual([]); + expect(protectedSweep.collected).toEqual([]); + expect(manager.getAgent(recent.id)?.lifecycle).toBe("idle"); + expect(manager.getAgent(protectedAgent.id)?.lifecycle).toBe("idle"); + expect(manager.getAgent(internal.id)?.lifecycle).toBe("idle"); + expect(manager.getAgent(running.id)?.lifecycle).toBe("running"); + expect(manager.getAgent(failed.id)?.lifecycle).toBe("error"); + } finally { + await Promise.all(manager.listAgents().map((agent) => manager.closeAgent(agent.id))).catch( + () => undefined, + ); + rmSync(workdir, { recursive: true, force: true }); + } +}); + +test("load waits for an in-flight collection close and creates only one resumed runtime", async () => { + const workdir = mkdtempSync(join(tmpdir(), "agent-manager-idle-close-race-")); + const storage = new AgentStorage(join(workdir, "agents"), logger); + const closeStarted = deferred(); + const closeAllowed = deferred(); + const client = new (class extends TestAgentClient { + resumeCount = 0; + + override async createSession(config: AgentSessionConfig): Promise { + return new (class extends TestAgentSession { + override async close(): Promise { + closeStarted.resolve(); + await closeAllowed.promise; + } + })(config); + } + + override async resumeSession( + handle: AgentPersistenceHandle, + config?: Partial, + ): Promise { + this.resumeCount += 1; + return super.resumeSession(handle, config); + } + })(); + const manager = new AgentManager({ clients: { codex: client }, registry: storage, logger }); + + try { + const created = await manager.createAgent( + { provider: "codex", cwd: workdir }, + "00000000-0000-4000-8000-000000000216", + { workspaceId: undefined }, + ); + const collection = manager.collectIdleAgents({ + cutoff: new Date(Date.now() + 1_000), + protectedAgentIds: new Set(), + }); + await closeStarted.promise; + const loads = Promise.all([ + ensureAgentLoaded(created.id, { agentManager: manager, agentStorage: storage, logger }), + ensureAgentLoaded(created.id, { agentManager: manager, agentStorage: storage, logger }), + ]); + + expect(client.resumeCount).toBe(0); + closeAllowed.resolve(); + const [first, second] = await loads; + await collection; + + expect(first.id).toBe(created.id); + expect(second.id).toBe(created.id); + expect(client.resumeCount).toBe(1); + } finally { + await manager.closeAgent("00000000-0000-4000-8000-000000000216").catch(() => undefined); + await storage.flush().catch(() => undefined); + rmSync(workdir, { recursive: true, force: true }); + } +}); + +test("provider close failure still persists and emits a resumable closed agent", async () => { + const workdir = mkdtempSync(join(tmpdir(), "agent-manager-close-failure-")); + const storage = new AgentStorage(join(workdir, "agents"), logger); + const client = new (class extends TestAgentClient { + override async createSession(config: AgentSessionConfig): Promise { + return new (class extends TestAgentSession { + override async close(): Promise { + throw new Error("provider cleanup failed"); + } + })(config); + } + })(); + const manager = new AgentManager({ clients: { codex: client }, registry: storage, logger }); + + try { + const created = await manager.createAgent( + { provider: "codex", cwd: workdir }, + "00000000-0000-4000-8000-000000000217", + { workspaceId: undefined }, + ); + const closed = waitForAgentLifecycle(manager, created.id, "closed"); + + const collection = await manager.collectIdleAgents({ + cutoff: new Date(Date.now() + 1_000), + protectedAgentIds: new Set(), + }); + await closed; + expect(collection.collected).toEqual([]); + expect(collection.failures).toHaveLength(1); + expect(collection.failures[0]).toMatchObject({ + agentId: created.id, + provider: "codex", + error: expect.objectContaining({ message: "provider cleanup failed" }), + }); + const stored = await storage.get(created.id); + expect(stored).toMatchObject({ lastStatus: "closed" }); + expect(stored?.archivedAt).toBeFalsy(); + + await expect( + ensureAgentLoaded(created.id, { agentManager: manager, agentStorage: storage, logger }), + ).resolves.toMatchObject({ id: created.id, lifecycle: "idle" }); + } finally { + await manager.closeAgent("00000000-0000-4000-8000-000000000217").catch(() => undefined); + await storage.flush().catch(() => undefined); + rmSync(workdir, { recursive: true, force: true }); + } +}); + test("hydrateTimeline keeps provider user_message items when no canonical user history exists", async () => { const workdir = mkdtempSync(join(tmpdir(), "agent-manager-history-keep-user-")); const storagePath = join(workdir, "agents"); diff --git a/packages/server/src/server/agent/agent-manager.ts b/packages/server/src/server/agent/agent-manager.ts index 158cdc811..6f0eb9cc0 100644 --- a/packages/server/src/server/agent/agent-manager.ts +++ b/packages/server/src/server/agent/agent-manager.ts @@ -391,6 +391,21 @@ export interface AgentMetricsSnapshot { }; } +export interface IdleAgentCollectionEntry { + agentId: string; + provider: AgentProvider; + sessionId?: string; +} + +export interface IdleAgentCollectionFailure extends IdleAgentCollectionEntry { + error: unknown; +} + +export interface IdleAgentCollectionResult { + collected: IdleAgentCollectionEntry[]; + failures: IdleAgentCollectionFailure[]; +} + type ActiveManagedAgent = | ManagedAgentInitializing | ManagedAgentIdle @@ -568,6 +583,7 @@ export class AgentManager { private readonly previousStatuses = new Map(); private readonly backgroundTasks = new Set>(); private readonly agentRegistrationTasks = new Set>(); + private readonly inFlightAgentCloses = new Map>(); private readonly agentStreamCoalescer: AgentStreamCoalescer; private mcpBaseUrl: string | null; private readonly mcpAuthToken: string | null; @@ -949,6 +965,19 @@ export class AgentManager { return agent ? { ...agent } : null; } + touchAgentActivity(id: string): ManagedAgent | null { + const agent = this.agents?.get(id); + if (!agent) { + return null; + } + this.touchUpdatedAt(agent); + return { ...agent }; + } + + async waitForAgentClose(agentId: string): Promise { + await this.inFlightAgentCloses?.get(agentId)?.catch(() => undefined); + } + getTimeline(id: string): AgentTimelineItem[] { this.requireAgent(id); return this.timelineStore.getItems(id); @@ -1004,6 +1033,7 @@ export class AgentManager { ): Promise { this.assertAcceptingAgentRegistrations(); const resolvedAgentId = validateAgentId(agentId ?? this.idFactory(), "createAgent"); + await this.deleteAgentState(resolvedAgentId); const { storedConfig, launchConfig } = await this.prepareSessionConfig(config, resolvedAgentId); this.requireEnabledProvider(storedConfig.provider); const client = await this.requireAvailableClient({ @@ -1088,7 +1118,10 @@ export class AgentManager { const launchContext = await this.buildLaunchContext(resolvedAgentId, client); const providerLaunchConfig = this.resolveProviderLaunchConfig(launchConfig, launchContext); const session = await client.resumeSession(handle, providerLaunchConfig, launchContext); - return this.registerSession(session, storedConfig, resolvedAgentId, options); + return this.registerSession(session, storedConfig, resolvedAgentId, { + ...options, + persistence: handle, + }); } importProviderSession(input: { @@ -1309,7 +1342,24 @@ export class AgentManager { } } - async closeAgent(agentId: string, options: { persistClosedState?: boolean } = {}): Promise { + closeAgent(agentId: string): Promise { + const existing = this.inFlightAgentCloses.get(agentId); + if (existing) { + return existing; + } + + const close = this.closeAgentRuntime(agentId); + this.inFlightAgentCloses.set(agentId, close); + const clearClose = () => { + if (this.inFlightAgentCloses.get(agentId) === close) { + this.inFlightAgentCloses.delete(agentId); + } + }; + void close.then(clearClose, clearClose); + return close; + } + + private async closeAgentRuntime(agentId: string): Promise { const agent = this.requireAgent(agentId); this.logger.trace( { @@ -1324,13 +1374,18 @@ export class AgentManager { "agent.manager.close.start", ); const closedAgent = this.prepareAgentForClosure(agent, "agent closed"); - await agent.session.close(); - this.timelineStore.delete(agentId); - for (const event of this.providerSubagents.deleteParent(agentId)) { - this.dispatch({ type: "provider_subagent", event }); + let closeError: unknown; + try { + await agent.session.close(); + } catch (error) { + closeError = error; } - if (options.persistClosedState !== false) { + + let persistError: unknown; + try { await this.persistSnapshot(closedAgent); + } catch (error) { + persistError = error; } this.emitClosedAgent(closedAgent, { persist: false }); this.logger.trace( @@ -1341,6 +1396,58 @@ export class AgentManager { }, "agent.manager.close.complete", ); + + if (closeError !== undefined) { + throw closeError; + } + if (persistError !== undefined) { + throw persistError; + } + } + + async collectIdleAgents(options: { + cutoff: Date; + protectedAgentIds: ReadonlySet; + }): Promise { + const result: IdleAgentCollectionResult = { collected: [], failures: [] }; + + for (const agent of Array.from(this.agents.values())) { + const current = this.agents.get(agent.id); + if (!current || !this.isIdleAgentCollectable(current, options)) { + continue; + } + + const entry: IdleAgentCollectionEntry = { + agentId: current.id, + provider: current.provider, + ...(current.persistence?.sessionId ? { sessionId: current.persistence.sessionId } : {}), + }; + try { + await this.closeAgent(current.id); + result.collected.push(entry); + } catch (error) { + result.failures.push({ ...entry, error }); + } + } + + return result; + } + + private isIdleAgentCollectable( + agent: LiveManagedAgent, + options: { cutoff: Date; protectedAgentIds: ReadonlySet }, + ): agent is ManagedAgentIdle { + return ( + agent.lifecycle === "idle" && + agent.updatedAt.getTime() <= options.cutoff.getTime() && + !agent.internal && + !options.protectedAgentIds.has(agent.id) && + agent.activeForegroundTurnId === null && + !this.runs.hasRun(agent.id) && + !agent.pendingReplacement && + agent.pendingPermissions.size === 0 && + agent.inFlightPermissionResponses.size === 0 + ); } async archiveAgent(agentId: string): Promise<{ archivedAt: string }> { @@ -1360,6 +1467,7 @@ export class AgentManager { const { archivedAt } = await this.markRecordArchived(stored); agent.updatedAt = new Date(archivedAt); await this.closeAgent(agentId); + this.discardRetainedAgentState(agentId); await this.cascadeArchiveChildren(agentId); @@ -1386,8 +1494,7 @@ export class AgentManager { if (this.agents.has(record.id)) { await this.archiveAgent(record.id); } else { - await this.markRecordArchived(record); - await this.cascadeArchiveChildren(record.id); + await this.archiveSnapshot(record.id, new Date().toISOString()); } } } @@ -1674,11 +1781,15 @@ export class AgentManager { if (this.agents.has(agentId)) { this.notifyAgentState(agentId); - } else if (!nextRecord.internal) { - this.dispatchArchivedStoredAgent(nextRecord); + } else { + this.discardRetainedAgentState(agentId); + if (!nextRecord.internal) { + this.dispatchArchivedStoredAgent(nextRecord); + } } await this.fireAgentArchived(agentId); + await this.cascadeArchiveChildren(agentId); return nextRecord; } @@ -2329,6 +2440,11 @@ export class AgentManager { await this.durableTimelineStore.deleteAgent(agentId); } + async deleteAgentState(agentId: string): Promise { + this.discardRetainedAgentState(agentId); + await this.deleteCommittedTimeline(agentId); + } + async getLastAssistantMessage(agentId: string): Promise { const agent = this.agents.get(agentId); if (!agent) { @@ -2638,6 +2754,7 @@ export class AgentManager { await this.refreshSessionState(managed, { emit: false }); this.assertAgentRegistrationActive(managed); managed.lifecycle = "idle"; + this.touchUpdatedAt(managed); await this.persistSnapshot(managed); this.assertAgentRegistrationActive(managed); this.emitState(managed, { persist: false }); @@ -2686,6 +2803,7 @@ export class AgentManager { | undefined; }): Promise<{ durableTimelineHasRows: boolean }> { const { agentId, now, options } = params; + const timelineAlreadyPrimed = this.timelineStore.has(agentId); const explicitTimelineSeed = buildExplicitTimelineSeedForRegister(now, options); const shouldSeedFromDurable = !explicitTimelineSeed && @@ -2695,7 +2813,8 @@ export class AgentManager { ? await this.loadCommittedTimelineSeed(agentId, now) : null; const durableTimelineHasRows = - durableTimelineSeed != null && (durableTimelineSeed.nextSeq ?? 1) > 1; + timelineAlreadyPrimed || + (durableTimelineSeed != null && (durableTimelineSeed.nextSeq ?? 1) > 1); const timelineSeed = explicitTimelineSeed ?? durableTimelineSeed; if (timelineSeed || !this.timelineStore.has(agentId)) { this.timelineStore.initialize(agentId, timelineSeed ?? { timestamp: now.toISOString() }); @@ -2803,9 +2922,23 @@ export class AgentManager { lifecycle: "closed", session: null, activeForegroundTurnId: null, + pendingPermissions: new Map(), + bufferedPermissionResolutions: new Map(), + inFlightPermissionResponses: new Set(), + pendingReplacement: false, + foregroundTurnWaiters: new Set(), + finalizedForegroundTurnIds: new Set(), + unsubscribeSession: null, }; } + private discardRetainedAgentState(agentId: string): void { + this.timelineStore.delete(agentId); + for (const event of this.providerSubagents.deleteParent(agentId)) { + this.dispatch({ type: "provider_subagent", event }); + } + } + private emitClosedAgent(agent: ManagedAgentClosed, options?: { persist?: boolean }): void { this.emitState(agent, options); } diff --git a/packages/server/src/server/agent/agent-response-loop.ts b/packages/server/src/server/agent/agent-response-loop.ts index d3a3de010..a76f82fda 100644 --- a/packages/server/src/server/agent/agent-response-loop.ts +++ b/packages/server/src/server/agent/agent-response-loop.ts @@ -382,6 +382,8 @@ export async function generateStructuredAgentResponse( await manager.closeAgent(agent.id); } catch { // ignore cleanup errors + } finally { + await manager.deleteAgentState(agent.id).catch(() => undefined); } } } diff --git a/packages/server/src/server/agent/agent-sdk-types.ts b/packages/server/src/server/agent/agent-sdk-types.ts index 16da95c05..3813e4eec 100644 --- a/packages/server/src/server/agent/agent-sdk-types.ts +++ b/packages/server/src/server/agent/agent-sdk-types.ts @@ -630,6 +630,7 @@ export interface AgentSession { ): Promise; describePersistence(): AgentPersistenceHandle | null; interrupt(): Promise; + /** Release live runtime resources without archiving or deleting the durable native session. */ close(): Promise; listCommands?(): Promise; setModel?(modelId: string | null): Promise; @@ -707,12 +708,12 @@ export interface AgentClient { isAvailable(): Promise; getDiagnostic?(): Promise<{ diagnostic: string }>; /** - * Archive a persisted session in the native provider (best-effort). + * Archive a durable native session (best-effort). Runtime release belongs to AgentSession.close(). * Called when Paseo archives an agent so the provider's own UI reflects the same state. */ archiveNativeSession?(handle: AgentPersistenceHandle): Promise; /** - * Unarchive a persisted session in the native provider. + * Unarchive a durable native session in the provider. * Called before Paseo clears its archived flag so provider resume can succeed. */ unarchiveNativeSession?(handle: AgentPersistenceHandle): Promise; diff --git a/packages/server/src/server/agent/providers/opencode-agent-commands.e2e.test.ts b/packages/server/src/server/agent/providers/opencode-agent-commands.e2e.test.ts index f712c6034..8dd3832cf 100644 --- a/packages/server/src/server/agent/providers/opencode-agent-commands.e2e.test.ts +++ b/packages/server/src/server/agent/providers/opencode-agent-commands.e2e.test.ts @@ -37,6 +37,28 @@ describe("opencode agent commands E2E", () => { } }, 60_000); + test("listing commands resumes an idle-collected agent", async () => { + const agent = await ctx.client.createAgent({ + ...getFullAccessConfig("opencode"), + cwd: "/tmp", + title: "Collected OpenCode Commands Test Agent", + }); + + const collection = await ctx.daemon.daemon.agentManager.collectIdleAgents({ + cutoff: new Date(Date.now() + 1_000), + protectedAgentIds: new Set(), + }); + expect(collection.failures).toEqual([]); + expect(collection.collected.map((entry) => entry.agentId)).toContain(agent.id); + expect(ctx.daemon.daemon.agentManager.getAgent(agent.id)).toBeNull(); + + const result = await ctx.client.listCommands({ agentId: agent.id }); + + expect(result.error).toBeNull(); + expect(result.commands.length).toBeGreaterThan(0); + expect(ctx.daemon.daemon.agentManager.getAgent(agent.id)?.id).toBe(agent.id); + }, 60_000); + test("sendMessage executes a slash command without arguments", async () => { const agent = await ctx.client.createAgent({ ...getFullAccessConfig("opencode"), diff --git a/packages/server/src/server/agent/providers/opencode-agent.test.ts b/packages/server/src/server/agent/providers/opencode-agent.test.ts index 829422ea7..e5ef7f681 100644 --- a/packages/server/src/server/agent/providers/opencode-agent.test.ts +++ b/packages/server/src/server/agent/providers/opencode-agent.test.ts @@ -194,7 +194,8 @@ describe("OpenCodeAgentClient adapter smoke tests", () => { test("creates a session with valid id and provider", async () => { const cwd = tmpCwd(); const runtime = new TestOpenCodeHarness(); - runtime.enqueueClient(new TestOpenCodeClient()); + const openCode = new TestOpenCodeClient(); + runtime.enqueueClient(openCode); const client = new OpenCodeAgentClient(logger, undefined, { serverManager: runtime, createClient: runtime.createClient, @@ -206,9 +207,49 @@ describe("OpenCodeAgentClient adapter smoke tests", () => { expect(session.provider).toBe("opencode"); await session.close(); + expect(openCode.calls.sessionAbort).toEqual([{ sessionID: "session-1", directory: cwd }]); + expect(openCode.calls.sessionUpdate).toEqual([]); rmSync(cwd, { recursive: true, force: true }); }, 60_000); + test("archives and unarchives the durable native session through client hooks", async () => { + const cwd = tmpCwd(); + const runtime = new TestOpenCodeHarness(); + const archiveClient = new TestOpenCodeClient(); + const unarchiveClient = new TestOpenCodeClient(); + runtime.enqueueClient(archiveClient); + runtime.enqueueClient(unarchiveClient); + const client = new OpenCodeAgentClient(logger, undefined, { + serverManager: runtime, + createClient: runtime.createClient, + }); + const handle = { + provider: "opencode" as const, + sessionId: "session-1", + metadata: { cwd }, + }; + + await client.archiveNativeSession(handle); + await client.unarchiveNativeSession(handle); + + expect(archiveClient.calls.sessionUpdate).toEqual([ + { + sessionID: "session-1", + directory: cwd, + time: { archived: expect.any(Number) }, + }, + ]); + expect(unarchiveClient.calls.sessionUpdate).toEqual([ + { + sessionID: "session-1", + directory: cwd, + time: { archived: null }, + }, + ]); + expect(runtime.acquisitions.every((acquisition) => acquisition.releaseCount === 1)).toBe(true); + rmSync(cwd, { recursive: true, force: true }); + }); + test("single turn completes with streaming deltas", async () => { const cwd = tmpCwd(); const runtime = new TestOpenCodeHarness(); @@ -713,65 +754,6 @@ describe("OpenCodeAgentClient adapter smoke tests", () => { }); describe("OpenCode adapter context-window normalization", () => { - test("close reconciliation aborts then archives upstream session", async () => { - const abort = vi.fn().mockResolvedValue({ data: true, error: undefined }); - const update = vi.fn().mockResolvedValue({ - data: { id: "session-1", time: { archived: Date.now() } }, - error: undefined, - }); - - await __openCodeInternals.reconcileOpenCodeSessionClose({ - client: { - session: { - abort, - update, - }, - } as never, - sessionId: "session-1", - directory: "/tmp/project", - logger: createTestLogger(), - }); - - expect(abort).toHaveBeenCalledWith({ - sessionID: "session-1", - directory: "/tmp/project", - }); - expect(update).toHaveBeenCalledTimes(1); - expect(update).toHaveBeenCalledWith({ - sessionID: "session-1", - directory: "/tmp/project", - time: { - archived: expect.any(Number), - }, - }); - }); - - test("close reconciliation still archives when abort returns an error", async () => { - const abort = vi.fn().mockResolvedValue({ - data: undefined, - error: { data: {}, errors: [], success: false }, - }); - const update = vi.fn().mockResolvedValue({ - data: { id: "session-1", time: { archived: Date.now() } }, - error: undefined, - }); - - await __openCodeInternals.reconcileOpenCodeSessionClose({ - client: { - session: { - abort, - update, - }, - } as never, - sessionId: "session-1", - directory: "/tmp/project", - logger: createTestLogger(), - }); - - expect(abort).toHaveBeenCalledTimes(1); - expect(update).toHaveBeenCalledTimes(1); - }); - test("builds OpenCode file parts for image prompt blocks", () => { expect( __openCodeInternals.buildOpenCodePromptParts([ @@ -2345,6 +2327,7 @@ describe("OpenCode persisted sessions", () => { describe("OpenCode provider subagent contract", () => { async function createAdoptedChildSession(): Promise<{ readonly runtime: TestOpenCodeHarness; + readonly provider: OpenCodeAgentClient; readonly parent: Awaited>; readonly child: Awaited>; readonly childClient: TestOpenCodeClient; @@ -2387,9 +2370,36 @@ describe("OpenCode provider subagent contract", () => { undefined, { env: { PASEO_AGENT_ID: "child-agent" } }, ); - return { runtime, parent, child, childClient }; + return { runtime, provider: client, parent, child, childClient }; } + test("archives an adopted child on the parent's registered OpenCode server", async () => { + const { runtime, provider, parent, child } = await createAdoptedChildSession(); + const archiveClient = new TestOpenCodeClient(); + runtime.enqueueClient(archiveClient); + + await provider.archiveNativeSession({ + provider: "opencode", + sessionId: "ses_child_external", + metadata: { cwd: "/workspace/repo" }, + }); + + expect(archiveClient.calls.sessionUpdate).toEqual([ + { + sessionID: "ses_child_external", + directory: "/workspace/repo", + time: { archived: expect.any(Number) }, + }, + ]); + expect(runtime.acquisitions.at(-1)).toEqual({ + kind: "existing", + url: runtime.server.url, + releaseCount: 1, + }); + await child.close(); + await parent.close(); + }); + test("resumes an adopted child on the parent's registered OpenCode server", async () => { const runtime = new TestOpenCodeHarness(); const parentClient = new TestOpenCodeClient(); diff --git a/packages/server/src/server/agent/providers/opencode-agent.ts b/packages/server/src/server/agent/providers/opencode-agent.ts index b70fdfa76..aabd1b627 100644 --- a/packages/server/src/server/agent/providers/opencode-agent.ts +++ b/packages/server/src/server/agent/providers/opencode-agent.ts @@ -431,7 +431,7 @@ function isOpenCodeNotFoundError(error: unknown): boolean { ); } -async function reconcileOpenCodeSessionClose(params: { +async function abortOpenCodeSession(params: { client: Pick; sessionId: string; directory: string; @@ -462,31 +462,6 @@ async function reconcileOpenCodeSessionClose(params: { "Failed to abort OpenCode session during close", ); } - - try { - const response = await client.session.update({ - sessionID: sessionId, - directory, - time: { archived: Date.now() }, - }); - if (response.error && !isOpenCodeNotFoundError(response.error)) { - logger.warn( - { - sessionId, - error: toDiagnosticErrorMessage(response.error), - }, - "Failed to archive OpenCode session during close", - ); - } - } catch (error) { - logger.warn( - { - sessionId, - error: toDiagnosticErrorMessage(error), - }, - "Failed to archive OpenCode session during close", - ); - } } function isOpenCodeHeadersTimeoutFailure(error: unknown): boolean { @@ -1254,7 +1229,6 @@ export const __openCodeInternals = { hasNormalizedOpenCodeUsage, mergeOpenCodeStepFinishUsage, parseOpenCodeModelLookupKey, - reconcileOpenCodeSessionClose, resolveOpenCodeModelLookupKeyFromAssistantMessage, resolveOpenCodeSelectedModelContextWindow, isSelectableOpenCodeAgent, @@ -1353,7 +1327,7 @@ export class OpenCodeAgentClient implements AgentClient { url, ); } catch (error) { - acquisition.release(); + await acquisition.release(); throw error; } } @@ -1407,7 +1381,7 @@ export class OpenCodeAgentClient implements AgentClient { registeredAcquisition !== null, ); } catch (error) { - acquisition.release(); + await acquisition.release(); throw error; } } @@ -1439,7 +1413,7 @@ export class OpenCodeAgentClient implements AgentClient { ]); return { models, modes }; } finally { - acquisition.release(); + await acquisition.release(); } } @@ -1455,7 +1429,7 @@ export class OpenCodeAgentClient implements AgentClient { try { return await listOpenCodeCommandsFromSdk(client, openCodeConfig.cwd); } finally { - acquisition.release(); + await acquisition.release(); } } @@ -1476,7 +1450,7 @@ export class OpenCodeAgentClient implements AgentClient { try { return await collectOpenCodeImportableSessionsFromSdk(client, options); } finally { - acquisition.release(); + await acquisition.release(); } } @@ -1512,7 +1486,57 @@ export class OpenCodeAgentClient implements AgentClient { }, }); } finally { - acquisition.release(); + await acquisition.release(); + } + } + + async archiveNativeSession(handle: AgentPersistenceHandle): Promise { + await this.setNativeSessionArchived(handle, Date.now()); + } + + async unarchiveNativeSession(handle: AgentPersistenceHandle): Promise { + await this.setNativeSessionArchived(handle, null); + } + + private async setNativeSessionArchived( + handle: AgentPersistenceHandle, + archivedAt: number | null, + ): Promise { + const metadata = (handle.metadata ?? {}) as Partial; + if (!metadata.cwd) { + throw new Error("OpenCode native archive update requires the original working directory"); + } + + const registeredServerUrl = getOpenCodeChildSessionServerUrl(handle.sessionId); + const acquisition = + (registeredServerUrl ? this.serverManager.acquireExisting(registeredServerUrl) : null) ?? + (await this.serverManager.acquireCurrent()); + const client = this.createOpenCodeClient({ + baseUrl: acquisition.server.url, + directory: metadata.cwd, + }); + try { + // OpenCode accepts null to clear the archive timestamp, but this SDK + // release's generated request type still exposes only number. + const updateSession = client.session.update.bind(client.session) as (parameters: { + sessionID: string; + directory?: string; + time?: { archived?: number | null }; + }) => ReturnType; + const response = readOpenCodeRecord( + await updateSession({ + sessionID: handle.sessionId, + directory: metadata.cwd, + time: { archived: archivedAt }, + }), + ); + if (response?.error) { + throw new Error( + `Failed to ${archivedAt === null ? "unarchive" : "archive"} OpenCode session: ${toDiagnosticErrorMessage(response.error)}`, + ); + } + } finally { + await acquisition.release(); } } @@ -2897,7 +2921,7 @@ class OpenCodeAgentSession implements AgentSession { private childHydrationCompleted = false; private readonly unrelatedSessionIds = new Set(); private selectedModelContextWindowMaxTokens: number | undefined; - private releaseServer: (() => void) | null; + private releaseServer: (() => Promise) | null; private eventStreamAbortController: AbortController | null = null; private eventStreamReady: Deferred | null = null; private eventStreamTask: Promise | null = null; @@ -2911,7 +2935,7 @@ class OpenCodeAgentSession implements AgentSession { sessionId: string, logger: Logger, modelContextWindowsByModelKey: ReadonlyMap = new Map(), - releaseServer?: () => void, + releaseServer?: () => Promise, persistSession = true, private readonly agentId?: string, private readonly serverUrl?: string, @@ -3891,7 +3915,7 @@ class OpenCodeAgentSession implements AgentSession { this.eventStreamReady = null; this.eventStreamTask = null; this.subscribers.clear(); - await reconcileOpenCodeSessionClose({ + await abortOpenCodeSession({ client: this.client, sessionId: this.sessionId, directory: this.config.cwd, @@ -3900,7 +3924,7 @@ class OpenCodeAgentSession implements AgentSession { await this.deleteProviderSessionIfEphemeral(); this.activeForegroundTurnId = null; } finally { - this.releaseServer?.(); + await this.releaseServer?.(); this.releaseServer = null; } } diff --git a/packages/server/src/server/agent/providers/opencode-server-manager.test.ts b/packages/server/src/server/agent/providers/opencode-server-manager.test.ts index 60c93c6ef..de094dba6 100644 --- a/packages/server/src/server/agent/providers/opencode-server-manager.test.ts +++ b/packages/server/src/server/agent/providers/opencode-server-manager.test.ts @@ -37,10 +37,10 @@ describe("OpenCodeServerManager generations", () => { expect(newAcquisition.server.url).toBe("http://127.0.0.1:4102"); expect(runtime.terminatedPorts).toEqual([]); - newAcquisition.release(); - oldAcquisition.release(); + await newAcquisition.release(); + await oldAcquisition.release(); - expect(runtime.terminatedPorts).toEqual([4101]); + expect(runtime.terminatedPorts).toEqual([4102, 4101]); }); test("new acquisitions after rotation use the new server", async () => { @@ -48,22 +48,22 @@ describe("OpenCodeServerManager generations", () => { const oldAcquisition = await manager.acquireCurrent(); const rotatedAcquisition = await manager.acquireNew(); - rotatedAcquisition.release(); - const nextAcquisition = await manager.acquireCurrent(); expect(nextAcquisition.server.url).toBe("http://127.0.0.1:4202"); expect(runtime.terminatedPorts).toEqual([]); - nextAcquisition.release(); - oldAcquisition.release(); + await rotatedAcquisition.release(); + expect(runtime.terminatedPorts).toEqual([]); + await nextAcquisition.release(); + await oldAcquisition.release(); }); test("concurrent new-server acquisitions share one fresh generation", async () => { const { manager, runtime } = createTestManager([4251, 4252, 4253]); const initialAcquisition = await manager.acquireCurrent(); - initialAcquisition.release(); + await initialAcquisition.release(); const [modelsAcquisition, modesAcquisition] = await Promise.all([ manager.acquireNew(), @@ -74,8 +74,8 @@ describe("OpenCodeServerManager generations", () => { expect(modesAcquisition.server.url).toBe("http://127.0.0.1:4252"); expect(runtime.launchedPorts).toEqual([4251, 4252]); - modesAcquisition.release(); - modelsAcquisition.release(); + await modesAcquisition.release(); + await modelsAcquisition.release(); }); test("release is idempotent", async () => { @@ -83,12 +83,12 @@ describe("OpenCodeServerManager generations", () => { const oldAcquisition = await manager.acquireCurrent(); const newAcquisition = await manager.acquireNew(); - newAcquisition.release(); + await newAcquisition.release(); - oldAcquisition.release(); - oldAcquisition.release(); + await oldAcquisition.release(); + await oldAcquisition.release(); - expect(runtime.terminatedPorts).toEqual([4301]); + expect(runtime.terminatedPorts).toEqual([4302, 4301]); }); test("shutdown kills current and retired servers", async () => { @@ -152,16 +152,16 @@ describe("OpenCodeServerManager generations", () => { const dedicatedStart = manager.acquireDedicated({ TEST_ENV: "custom" }); await runtime.settle(); - currentAcquisition.release(); - expect(runtime.terminatedPorts).toEqual([]); + await currentAcquisition.release(); + expect(runtime.terminatedPorts).toEqual([4473]); runtime.processForPort(4474).announceListening(); const dedicatedAcquisition = await dedicatedStart; expect(dedicatedAcquisition.server.url).toBe("http://127.0.0.1:4474"); - dedicatedAcquisition.release(); - expect(runtime.terminatedPorts).toEqual([4474]); + await dedicatedAcquisition.release(); + expect(runtime.terminatedPorts).toEqual([4473, 4474]); }); test("acquireExisting keeps a retired dedicated server alive until every reference releases", async () => { @@ -172,10 +172,10 @@ describe("OpenCodeServerManager generations", () => { expect(existingAcquisition?.server.url).toBe("http://127.0.0.1:4475"); - dedicatedAcquisition.release(); + await dedicatedAcquisition.release(); expect(runtime.terminatedPorts).toEqual([]); - existingAcquisition?.release(); + await existingAcquisition?.release(); expect(runtime.terminatedPorts).toEqual([4475]); }); @@ -187,7 +187,7 @@ describe("OpenCodeServerManager generations", () => { expect(manager.acquireExisting("http://127.0.0.1:9999")).toBe(null); - acquisition.release(); + await acquisition.release(); expect(runtime.terminatedPorts).toEqual([4476]); expect(manager.acquireExisting(url)).toBe(null); }); @@ -197,12 +197,26 @@ describe("OpenCodeServerManager generations", () => { const firstAcquisition = await manager.acquireCurrent(); const secondAcquisition = await manager.acquireNew(); - secondAcquisition.release(); + await secondAcquisition.release(); const thirdAcquisition = await manager.acquireNew(); - thirdAcquisition.release(); - firstAcquisition.release(); + await thirdAcquisition.release(); + await firstAcquisition.release(); - expect(runtime.terminatedPorts).toEqual([4502, 4501]); + expect(runtime.terminatedPorts).toEqual([4502, 4503, 4501]); + }); + + test("final release detaches the terminating generation before a concurrent acquire", async () => { + const { manager, runtime } = createTestManager([4551, 4552]); + + const first = await manager.acquireCurrent(); + const release = first.release(); + const next = await manager.acquireCurrent(); + + await release; + expect(next.server.url).toBe("http://127.0.0.1:4552"); + expect(runtime.terminatedPorts).toEqual([4551]); + + await next.release(); }); }); @@ -258,7 +272,7 @@ describe("OpenCodeServerManager managed process ledger", () => { }), ]); - acquisition.release(); + await acquisition.release(); await manager.shutdown(); } finally { rmSync(tempDir, { recursive: true, force: true }); diff --git a/packages/server/src/server/agent/providers/opencode/server-manager.ts b/packages/server/src/server/agent/providers/opencode/server-manager.ts index 3fbb89ad2..d765d3d27 100644 --- a/packages/server/src/server/agent/providers/opencode/server-manager.ts +++ b/packages/server/src/server/agent/providers/opencode/server-manager.ts @@ -21,11 +21,10 @@ const OPENCODE_SERVER_FORCE_SHUTDOWN_TIMEOUT_MS = 1_000; export interface OpenCodeServerAcquisition { server: { port: number; url: string }; - release: () => void; + release: () => Promise; } export interface OpenCodeServerManagerLike { - ensureRunning(): Promise<{ port: number; url: string }>; acquireCurrent(): Promise; acquireNew(): Promise; acquireDedicated(env: Record): Promise; @@ -135,12 +134,6 @@ export class OpenCodeServerManager implements OpenCodeServerManagerLike { process.on("SIGINT", cleanup); } - async ensureRunning(): Promise<{ port: number; url: string }> { - const acquisition = await this.acquireCurrent(); - acquisition.release(); - return acquisition.server; - } - async acquireCurrent(): Promise { const server = await this.getCurrentServer(); return this.acquireServer(server); @@ -160,7 +153,7 @@ export class OpenCodeServerManager implements OpenCodeServerManagerLike { await server.ready; return acquisition; } catch (error) { - acquisition.release(); + await acquisition.release(); throw error; } } @@ -188,20 +181,37 @@ export class OpenCodeServerManager implements OpenCodeServerManagerLike { private acquireServer(server: OpenCodeServerGeneration): OpenCodeServerAcquisition { server.refCount += 1; - let released = false; + let releasePromise: Promise | null = null; return { server: { port: server.port, url: server.url }, - release: () => { - if (released) { - return; + release: async () => { + if (releasePromise) { + return releasePromise; } - released = true; - server.refCount -= 1; - this.cleanupRetiredServers(); + releasePromise = this.releaseServer(server); + return releasePromise; }, }; } + private async releaseServer(server: OpenCodeServerGeneration): Promise { + server.refCount = Math.max(0, server.refCount - 1); + if (server.refCount > 0) { + return; + } + + if (this.currentServer === server) { + this.currentServer = null; + server.retired = true; + } + if (!server.retired) { + return; + } + + this.retiredServers.delete(server); + await this.killServer(server); + } + private async getNewServer(): Promise { if (this.newServerPromise) { return this.newServerPromise; @@ -261,14 +271,14 @@ export class OpenCodeServerManager implements OpenCodeServerManagerLike { existing.retired = true; this.retiredServers.add(existing); this.currentServer = null; - this.cleanupRetiredServers(); + await this.cleanupRetiredServers(); } if (this.startPromise) { const pending = await this.startPromise; pending.retired = true; this.retiredServers.add(pending); this.currentServer = null; - this.cleanupRetiredServers(); + await this.cleanupRetiredServers(); } } @@ -417,13 +427,15 @@ export class OpenCodeServerManager implements OpenCodeServerManagerLike { this.retiredServers.clear(); } - private cleanupRetiredServers(): void { + private async cleanupRetiredServers(): Promise { + const cleanup: Promise[] = []; for (const server of Array.from(this.retiredServers)) { if (server.refCount === 0) { this.retiredServers.delete(server); - void this.killServer(server); + cleanup.push(this.killServer(server)); } } + await Promise.all(cleanup); } private async killServer(server: OpenCodeServerGeneration): Promise { diff --git a/packages/server/src/server/agent/providers/opencode/test-server-manager.ts b/packages/server/src/server/agent/providers/opencode/test-server-manager.ts index f13a2eb95..e4cc8a65e 100644 --- a/packages/server/src/server/agent/providers/opencode/test-server-manager.ts +++ b/packages/server/src/server/agent/providers/opencode/test-server-manager.ts @@ -10,12 +10,6 @@ export interface TestOpenCodeServerAcquisition { export class TestOpenCodeServerManager implements OpenCodeServerManagerLike { readonly acquisitions: TestOpenCodeServerAcquisition[] = []; readonly server = { port: 1234, url: "http://127.0.0.1:1234" }; - ensureRunningCount = 0; - - async ensureRunning(): Promise<{ port: number; url: string }> { - this.ensureRunningCount += 1; - return this.server; - } async acquireCurrent(): Promise { return this.recordAcquisition({ kind: "current" }); @@ -47,7 +41,7 @@ export class TestOpenCodeServerManager implements OpenCodeServerManagerLike { this.acquisitions.push(acquisition); return { server: this.server, - release: () => { + release: async () => { acquisition.released = true; }, }; diff --git a/packages/server/src/server/agent/providers/opencode/test-utils/test-opencode-harness.ts b/packages/server/src/server/agent/providers/opencode/test-utils/test-opencode-harness.ts index 754b1b86d..3c3ac4dc7 100644 --- a/packages/server/src/server/agent/providers/opencode/test-utils/test-opencode-harness.ts +++ b/packages/server/src/server/agent/providers/opencode/test-utils/test-opencode-harness.ts @@ -53,16 +53,12 @@ export class TestOpenCodeHarness implements OpenCodeServerManagerLike { this.acquisitions.push(acquisition); return { server: this.server, - release: () => { + release: async () => { acquisition.releaseCount += 1; }, }; } - async ensureRunning(): Promise<{ port: number; url: string }> { - return this.server; - } - readonly createClient = (options: { baseUrl: string; directory: string }): OpencodeClient => { this.clientCreations.push(options); const client = this.clients.shift() ?? new TestOpenCodeClient(); diff --git a/packages/server/src/server/bootstrap.ts b/packages/server/src/server/bootstrap.ts index c9d27c881..4d3d48929 100644 --- a/packages/server/src/server/bootstrap.ts +++ b/packages/server/src/server/bootstrap.ts @@ -217,6 +217,8 @@ import { DaemonExecutions } from "./hub/daemon-executions.js"; const MAX_MCP_DEBUG_BATCH_ITEMS = 10; const REDACTED_LOG_VALUE = "[redacted]"; +const IDLE_AGENT_RUNTIME_TTL_MS = 2 * 60 * 1000; +const IDLE_AGENT_RUNTIME_SWEEP_INTERVAL_MS = 15 * 1000; const DOWNLOAD_OPEN_FLAGS = process.platform === "win32" ? constants.O_RDONLY : constants.O_RDONLY | constants.O_NOFOLLOW; @@ -1175,6 +1177,39 @@ export async function createPaseoDaemon( archiveWorkspace: archiveScheduleWorkspaceExternal, }); await scheduleService.start(); + let inFlightIdleAgentCollection: Promise | null = null; + const collectIdleAgentRuntimes = async () => { + const protectedAgentIds = await scheduleService.listActiveAgentTargetIds(); + const cutoff = new Date(Date.now() - IDLE_AGENT_RUNTIME_TTL_MS); + const result = await agentManager.collectIdleAgents({ cutoff, protectedAgentIds }); + for (const collected of result.collected) { + logger.info(collected, "Collected idle agent runtime"); + } + for (const failure of result.failures) { + const { error, ...context } = failure; + logger.warn({ ...context, err: error }, "Failed to collect idle agent runtime"); + } + }; + const runIdleAgentCollection = () => { + if (inFlightIdleAgentCollection) { + return; + } + const collection = collectIdleAgentRuntimes() + .catch((error) => { + logger.warn({ err: error }, "Idle agent runtime sweep failed"); + }) + .finally(() => { + if (inFlightIdleAgentCollection === collection) { + inFlightIdleAgentCollection = null; + } + }); + inFlightIdleAgentCollection = collection; + }; + const idleAgentCollectionTimer = setInterval( + runIdleAgentCollection, + IDLE_AGENT_RUNTIME_SWEEP_INTERVAL_MS, + ); + idleAgentCollectionTimer.unref(); agentManager.setAgentArchivedCallback(async (agentId) => { try { await scheduleService.completeForAgent(agentId); @@ -1564,6 +1599,8 @@ export async function createPaseoDaemon( await hubRelationships.stop(); workspaceReconciliation.dispose(); scriptHealthMonitor.stop(); + clearInterval(idleAgentCollectionTimer); + await inFlightIdleAgentCollection; // Freeze both ingress and registration before taking the agent closure snapshot. wsServer?.prepareForShutdown(); agentManager.prepareForShutdown(); diff --git a/packages/server/src/server/hub/daemon-executions.ts b/packages/server/src/server/hub/daemon-executions.ts index 83b4d0693..712a7ecc2 100644 --- a/packages/server/src/server/hub/daemon-executions.ts +++ b/packages/server/src/server/hub/daemon-executions.ts @@ -177,7 +177,11 @@ export class DaemonExecutions implements HubExecutionAgents { try { await autoArchiveRegistration.cancel(); if (createdAgentId && this.agentManager.getAgent(createdAgentId)) { - await this.agentManager.closeAgent(createdAgentId); + try { + await this.agentManager.closeAgent(createdAgentId); + } finally { + await this.agentManager.deleteAgentState(createdAgentId); + } } } finally { try { diff --git a/packages/server/src/server/loop-service.ts b/packages/server/src/server/loop-service.ts index 99e5f8ff3..43474796c 100644 --- a/packages/server/src/server/loop-service.ts +++ b/packages/server/src/server/loop-service.ts @@ -241,6 +241,7 @@ type LoopAgentManager = Pick< | "archiveAgent" | "cancelAgentRun" | "closeAgent" + | "deleteAgentState" | "getAgent" | "runAgent" | "subscribe" @@ -570,7 +571,7 @@ export class LoopService { await this.options.agentManager.archiveAgent(agentId); return; } - await this.options.agentManager.closeAgent(agentId); + await this.closeInternalAgent(agentId); } catch (error) { if (!isUnknownLoopAgentError(error, agentId)) { throw error; @@ -578,6 +579,14 @@ export class LoopService { } } + private async closeInternalAgent(agentId: string): Promise { + try { + await this.options.agentManager.closeAgent(agentId); + } finally { + await this.options.agentManager.deleteAgentState(agentId); + } + } + private async executeLoop(loopId: string, signal: AbortSignal): Promise { const loop = this.requireLoop(loopId); const deadline = loop.maxTimeMs ? Date.now() + loop.maxTimeMs : null; @@ -775,7 +784,7 @@ export class LoopService { if (loop.archive) { await this.options.agentManager.archiveAgent(agent.id); } else { - await this.options.agentManager.closeAgent(agent.id); + await this.closeInternalAgent(agent.id); } } catch { // Ignore cleanup errors for internal loop workers. @@ -917,7 +926,7 @@ export class LoopService { if (loop.archive) { await this.options.agentManager.archiveAgent(verifierAgent.id); } else { - await this.options.agentManager.closeAgent(verifierAgent.id); + await this.closeInternalAgent(verifierAgent.id); } } catch { // Ignore cleanup errors for internal loop verifiers. diff --git a/packages/server/src/server/schedule/service.test.ts b/packages/server/src/server/schedule/service.test.ts index a7fad87be..b713641e0 100644 --- a/packages/server/src/server/schedule/service.test.ts +++ b/packages/server/src/server/schedule/service.test.ts @@ -376,6 +376,47 @@ describe("ScheduleService", () => { expect(resumed.nextRunAt).toBe("2026-01-01T00:04:00.000Z"); }); + test("lists only active schedules that target existing agents", async () => { + const service = createScheduleService({ + paseoHome: tempDir, + logger: createTestLogger(), + agentManager: new AgentManager({ logger: createTestLogger() }), + agentStorage, + providerSnapshotManager: NO_UNATTENDED_SCHEDULE_POLICY, + now: () => now, + runner: async () => ({ agentId: null, output: "ok" }), + }); + const activeAgentId = "00000000-0000-4000-8000-000000000201"; + const pausedAgentId = "00000000-0000-4000-8000-000000000202"; + const completedAgentId = "00000000-0000-4000-8000-000000000203"; + const cadence = { type: "every" as const, everyMs: 60_000 }; + + await service.create({ + prompt: "Keep active agent resident", + cadence, + target: { type: "agent", agentId: activeAgentId }, + }); + const paused = await service.create({ + prompt: "Paused heartbeat", + cadence, + target: { type: "agent", agentId: pausedAgentId }, + }); + await service.pause(paused.id); + await service.create({ + prompt: "Completed heartbeat", + cadence, + target: { type: "agent", agentId: completedAgentId }, + }); + await service.completeForAgent(completedAgentId); + await service.create({ + prompt: "Fresh agent each run", + cadence, + target: { type: "new-agent", config: { provider: "claude", cwd: tempDir } }, + }); + + await expect(service.listActiveAgentTargetIds()).resolves.toEqual(new Set([activeAgentId])); + }); + test("completes schedules when max runs is reached", async () => { const service = createScheduleService({ paseoHome: tempDir, diff --git a/packages/server/src/server/schedule/service.ts b/packages/server/src/server/schedule/service.ts index d55ecb464..6bcaed724 100644 --- a/packages/server/src/server/schedule/service.ts +++ b/packages/server/src/server/schedule/service.ts @@ -204,7 +204,9 @@ type ScheduleAgentManager = Pick< | "hydrateTimelineFromProvider" | "resumeAgentFromPersistence" | "runAgent" + | "touchAgentActivity" | "waitForAgentEvent" + | "waitForAgentClose" >; interface ScheduleWorkspaceCreateInput { @@ -368,6 +370,17 @@ export class ScheduleService { return this.store.list(); } + async listActiveAgentTargetIds(): Promise> { + const schedules = await this.store.list(); + const agentIds = new Set(); + for (const schedule of schedules) { + if (schedule.status === "active" && schedule.target.type === "agent") { + agentIds.add(schedule.target.agentId); + } + } + return agentIds; + } + async inspect(id: string): Promise { const schedule = await this.store.get(id); if (!schedule) { diff --git a/packages/server/src/server/session.ts b/packages/server/src/server/session.ts index 7d7ad185a..17510db81 100644 --- a/packages/server/src/server/session.ts +++ b/packages/server/src/server/session.ts @@ -2164,7 +2164,7 @@ export class Session { try { await this.agentStorage.remove(agentId); - await this.agentManager.deleteCommittedTimeline(agentId); + await this.agentManager.deleteAgentState(agentId); } catch (error) { this.sessionLogger.error({ err: error, agentId }, `Failed to fully delete agent ${agentId}`); } @@ -3350,6 +3350,15 @@ export class Session { const agentIds = Array.isArray(agentId) ? agentId : [agentId]; try { + await Promise.all( + agentIds.map((id) => + ensureAgentLoaded(id, { + agentManager: this.agentManager, + agentStorage: this.agentStorage, + logger: this.sessionLogger, + }), + ), + ); await Promise.all(agentIds.map((id) => this.agentManager.clearAgentAttention(id))); if (requestId) { const agents = ( @@ -3436,8 +3445,16 @@ export class Session { ); try { - const agents = this.agentManager.listAgents(); - const agent = agents.find((a) => a.id === agentId); + const existing = this.agentManager.getAgent(agentId); + const stored = existing ? null : await this.agentStorage.get(agentId); + const agent = + existing || (stored && !stored.archivedAt) + ? await ensureAgentLoaded(agentId, { + agentManager: this.agentManager, + agentStorage: this.agentStorage, + logger: this.sessionLogger, + }) + : null; if (agent?.session?.listCommands) { const commands = await agent.session.listCommands(); @@ -5847,6 +5864,11 @@ export class Session { msg: Extract, ): Promise { try { + await ensureAgentLoaded(msg.parentAgentId, { + agentManager: this.agentManager, + agentStorage: this.agentStorage, + logger: this.sessionLogger, + }); this.emit({ type: "agent.provider_subagents.list.response", payload: { @@ -5874,6 +5896,11 @@ export class Session { ): Promise { const direction: AgentTimelineFetchDirection = msg.direction ?? (msg.cursor ? "after" : "tail"); try { + await ensureAgentLoaded(msg.parentAgentId, { + agentManager: this.agentManager, + agentStorage: this.agentStorage, + logger: this.sessionLogger, + }); const descriptor = this.agentManager.getProviderSubagent(msg.parentAgentId, msg.subagentId); if (!descriptor) { throw new Error("Provider subagent not found");