From 8cc2ae0ba434ed4c500a7580fcab94f10cd53523 Mon Sep 17 00:00:00 2001 From: Mohamed Boudra Date: Sat, 18 Jul 2026 22:15:27 +0200 Subject: [PATCH] Resume collected agents before pane actions (#2209) * fix: resume collected agents before pane actions * fix(server): keep archived agents closed during pane actions An archive could win while a collected agent was resuming, then lose when the provider runtime registered. Recheck persisted archive state after registration and close the resumed runtime before any mutation runs. * fix(server): fence archived agents after shared resume Protected pane actions could join a resume started by an ordinary loader and skip the archive fence. Recheck persisted lifecycle state for every protected caller so archive always wins before mutation. --- .../server/src/server/agent/agent-loading.ts | 23 +++ .../src/server/agent/agent-manager.test.ts | 133 +++++++++++++++++- packages/server/src/server/session.test.ts | 37 +++-- packages/server/src/server/session.ts | 13 +- .../agent-config/agent-config-session.test.ts | 45 ++++++ .../agent-config/agent-config-session.ts | 8 +- 6 files changed, 243 insertions(+), 16 deletions(-) diff --git a/packages/server/src/server/agent/agent-loading.ts b/packages/server/src/server/agent/agent-loading.ts index 01dd9264b..e0e806b01 100644 --- a/packages/server/src/server/agent/agent-loading.ts +++ b/packages/server/src/server/agent/agent-loading.ts @@ -30,6 +30,29 @@ export interface EnsureAgentLoadedDeps { logger: Logger; } +export async function ensureUnarchivedAgentLoaded( + agentId: string, + deps: EnsureAgentLoadedDeps & { + agentManager: AgentLoaderManager & Pick; + }, +): Promise { + const record = await deps.agentStorage.get(agentId); + if (record?.archivedAt) { + throw new Error(`Agent is archived: ${agentId}`); + } + + const agent = await ensureAgentLoaded(agentId, deps); + const latestRecord = await deps.agentStorage.get(agentId); + if (latestRecord?.archivedAt) { + await deps.agentManager.closeAgent(agentId).catch((error: unknown) => { + deps.logger.warn({ err: error, agentId }, "Failed to close concurrently archived agent"); + }); + throw new Error(`Agent is archived: ${agentId}`); + } + + return agent; +} + export async function ensureAgentLoaded( agentId: string, deps: EnsureAgentLoadedDeps, diff --git a/packages/server/src/server/agent/agent-manager.test.ts b/packages/server/src/server/agent/agent-manager.test.ts index 12dea4e47..161a2b416 100644 --- a/packages/server/src/server/agent/agent-manager.test.ts +++ b/packages/server/src/server/agent/agent-manager.test.ts @@ -17,7 +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 { ensureAgentLoaded, ensureUnarchivedAgentLoaded } from "./agent-loading.js"; import type { StoredAgentRecord } from "./agent-storage.js"; import type { AgentClient, @@ -7300,6 +7300,137 @@ test("archiving an idle-collected parent still cascades to its managed children" } }); +test("ensureUnarchivedAgentLoaded does not resume an archived agent", async () => { + const workdir = mkdtempSync(join(tmpdir(), "agent-manager-archived-load-")); + const storage = new AgentStorage(join(workdir, "agents"), logger); + const manager = new AgentManager({ + clients: { codex: new TestAgentClient() }, + registry: storage, + logger, + }); + + try { + const agent = await manager.createAgent({ provider: "codex", cwd: workdir }, undefined, { + workspaceId: undefined, + }); + await manager.collectIdleAgents({ + cutoff: new Date(Date.now() + 1_000), + protectedAgentIds: new Set(), + }); + await manager.archiveSnapshot(agent.id, new Date().toISOString()); + + await expect( + ensureUnarchivedAgentLoaded(agent.id, { + agentManager: manager, + agentStorage: storage, + logger, + }), + ).rejects.toThrow(`Agent is archived: ${agent.id}`); + expect(manager.getAgent(agent.id)).toBeNull(); + } finally { + await storage.flush().catch(() => undefined); + rmSync(workdir, { recursive: true, force: true }); + } +}); + +test("ensureUnarchivedAgentLoaded closes a runtime archived while it resumes", async () => { + const workdir = mkdtempSync(join(tmpdir(), "agent-manager-archived-resume-race-")); + const storage = new AgentStorage(join(workdir, "agents"), logger); + const resumeStarted = deferred(); + const resumeAllowed = deferred(); + const client = new (class extends TestAgentClient { + override async resumeSession( + handle: AgentPersistenceHandle, + config?: Partial, + launchContext?: AgentLaunchContext, + ): Promise { + resumeStarted.resolve(); + await resumeAllowed.promise; + return super.resumeSession(handle, config, launchContext); + } + })(); + const manager = new AgentManager({ clients: { codex: client }, registry: storage, logger }); + + try { + const agent = await manager.createAgent({ provider: "codex", cwd: workdir }, undefined, { + workspaceId: undefined, + }); + await manager.collectIdleAgents({ + cutoff: new Date(Date.now() + 1_000), + protectedAgentIds: new Set(), + }); + + const load = ensureUnarchivedAgentLoaded(agent.id, { + agentManager: manager, + agentStorage: storage, + logger, + }); + await resumeStarted.promise; + await manager.archiveSnapshot(agent.id, new Date().toISOString()); + resumeAllowed.resolve(); + + await expect(load).rejects.toThrow(`Agent is archived: ${agent.id}`); + expect(manager.getAgent(agent.id)).toBeNull(); + expect((await storage.get(agent.id))?.archivedAt).toEqual(expect.any(String)); + } finally { + resumeAllowed.resolve(); + await storage.flush().catch(() => undefined); + rmSync(workdir, { recursive: true, force: true }); + } +}); + +test("ensureUnarchivedAgentLoaded fences an archived agent after joining a shared resume", async () => { + const workdir = mkdtempSync(join(tmpdir(), "agent-manager-archived-shared-resume-")); + const storage = new AgentStorage(join(workdir, "agents"), logger); + const resumeStarted = deferred(); + const resumeAllowed = deferred(); + const client = new (class extends TestAgentClient { + override async resumeSession( + handle: AgentPersistenceHandle, + config?: Partial, + launchContext?: AgentLaunchContext, + ): Promise { + resumeStarted.resolve(); + await resumeAllowed.promise; + return super.resumeSession(handle, config, launchContext); + } + })(); + const manager = new AgentManager({ clients: { codex: client }, registry: storage, logger }); + + try { + const agent = await manager.createAgent({ provider: "codex", cwd: workdir }, undefined, { + workspaceId: undefined, + }); + await manager.collectIdleAgents({ + cutoff: new Date(Date.now() + 1_000), + protectedAgentIds: new Set(), + }); + + const sharedLoad = ensureAgentLoaded(agent.id, { + agentManager: manager, + agentStorage: storage, + logger, + }); + await resumeStarted.promise; + const protectedLoad = ensureUnarchivedAgentLoaded(agent.id, { + agentManager: manager, + agentStorage: storage, + logger, + }); + await manager.archiveSnapshot(agent.id, new Date().toISOString()); + resumeAllowed.resolve(); + + await sharedLoad; + await expect(protectedLoad).rejects.toThrow(`Agent is archived: ${agent.id}`); + expect(manager.getAgent(agent.id)).toBeNull(); + expect((await storage.get(agent.id))?.archivedAt).toEqual(expect.any(String)); + } finally { + resumeAllowed.resolve(); + 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 { diff --git a/packages/server/src/server/session.test.ts b/packages/server/src/server/session.test.ts index 4db2bc7a6..183f02633 100644 --- a/packages/server/src/server/session.test.ts +++ b/packages/server/src/server/session.test.ts @@ -370,6 +370,7 @@ function createSessionForTest(options: SessionForTestOptions = {}): Session { ...options.agentManager, }), agentStorage: asAgentStorage({ + get: vi.fn().mockResolvedValue(undefined), list: vi.fn().mockResolvedValue([]), ...options.agentStorage, }), @@ -4899,12 +4900,22 @@ test("sends project updates only to capable sockets in a retained session", () = }); describe("agent config setters", () => { + function liveAgentManager(overrides: { [K in keyof SessionOptions["agentManager"]]?: unknown }): { + [K in keyof SessionOptions["agentManager"]]?: unknown; + } { + return { + waitForAgentClose: vi.fn().mockResolvedValue(undefined), + touchAgentActivity: vi.fn(() => ({ id: "agent-1" })), + ...overrides, + }; + } + test("set_agent_mode_request: success emits accepted response carrying the notice", async () => { const messages: SessionOutboundMessage[] = []; const notice = { type: "info", message: "Switched to plan mode" } as const; const session = createSessionForTest({ messages, - agentManager: { setAgentMode: vi.fn().mockResolvedValue(notice) }, + agentManager: liveAgentManager({ setAgentMode: vi.fn().mockResolvedValue(notice) }), }); await session.handleMessage({ @@ -4932,7 +4943,9 @@ describe("agent config setters", () => { const messages: SessionOutboundMessage[] = []; const session = createSessionForTest({ messages, - agentManager: { setAgentMode: vi.fn().mockRejectedValue(new Error("mode boom")) }, + agentManager: liveAgentManager({ + setAgentMode: vi.fn().mockRejectedValue(new Error("mode boom")), + }), }); await session.handleMessage({ @@ -4967,7 +4980,7 @@ describe("agent config setters", () => { const messages: SessionOutboundMessage[] = []; const session = createSessionForTest({ messages, - agentManager: { setAgentModel: vi.fn().mockResolvedValue(undefined) }, + agentManager: liveAgentManager({ setAgentModel: vi.fn().mockResolvedValue(undefined) }), }); await session.handleMessage({ @@ -4989,7 +5002,9 @@ describe("agent config setters", () => { const messages: SessionOutboundMessage[] = []; const session = createSessionForTest({ messages, - agentManager: { setAgentModel: vi.fn().mockRejectedValue(new Error("model boom")) }, + agentManager: liveAgentManager({ + setAgentModel: vi.fn().mockRejectedValue(new Error("model boom")), + }), }); await session.handleMessage({ @@ -5024,7 +5039,7 @@ describe("agent config setters", () => { const messages: SessionOutboundMessage[] = []; const session = createSessionForTest({ messages, - agentManager: { setAgentFeature: vi.fn().mockResolvedValue(undefined) }, + agentManager: liveAgentManager({ setAgentFeature: vi.fn().mockResolvedValue(undefined) }), }); await session.handleMessage({ @@ -5047,7 +5062,9 @@ describe("agent config setters", () => { const messages: SessionOutboundMessage[] = []; const session = createSessionForTest({ messages, - agentManager: { setAgentFeature: vi.fn().mockRejectedValue(new Error("feature boom")) }, + agentManager: liveAgentManager({ + setAgentFeature: vi.fn().mockRejectedValue(new Error("feature boom")), + }), }); await session.handleMessage({ @@ -5084,7 +5101,9 @@ describe("agent config setters", () => { const notice = { type: "warning", message: "Thinking budget reduced" } as const; const session = createSessionForTest({ messages, - agentManager: { setAgentThinkingOption: vi.fn().mockResolvedValue(notice) }, + agentManager: liveAgentManager({ + setAgentThinkingOption: vi.fn().mockResolvedValue(notice), + }), }); await session.handleMessage({ @@ -5112,9 +5131,9 @@ describe("agent config setters", () => { const messages: SessionOutboundMessage[] = []; const session = createSessionForTest({ messages, - agentManager: { + agentManager: liveAgentManager({ setAgentThinkingOption: vi.fn().mockRejectedValue(new Error("thinking boom")), - }, + }), }); await session.handleMessage({ diff --git a/packages/server/src/server/session.ts b/packages/server/src/server/session.ts index 17510db81..b974559cb 100644 --- a/packages/server/src/server/session.ts +++ b/packages/server/src/server/session.ts @@ -36,7 +36,7 @@ import { isStoredAgentProviderAvailable, toAgentPersistenceHandle, } from "./persistence-hooks.js"; -import { ensureAgentLoaded } from "./agent/agent-loading.js"; +import { ensureAgentLoaded, ensureUnarchivedAgentLoaded } from "./agent/agent-loading.js"; import { formatSystemNotificationPrompt, sendPromptToAgent, @@ -816,6 +816,13 @@ export class Session { emit: (msg) => this.emit(msg), }, operations: { + ensureLoaded: async (agentId) => { + await ensureUnarchivedAgentLoaded(agentId, { + agentManager, + agentStorage, + logger: this.sessionLogger, + }); + }, setMode: async (agentId, modeId) => (await setAgentModeCommand({ agentManager }, { agentId, modeId })).notice, setModel: (agentId, modelId) => agentManager.setAgentModel(agentId, modelId), @@ -5864,7 +5871,7 @@ export class Session { msg: Extract, ): Promise { try { - await ensureAgentLoaded(msg.parentAgentId, { + await ensureUnarchivedAgentLoaded(msg.parentAgentId, { agentManager: this.agentManager, agentStorage: this.agentStorage, logger: this.sessionLogger, @@ -5896,7 +5903,7 @@ export class Session { ): Promise { const direction: AgentTimelineFetchDirection = msg.direction ?? (msg.cursor ? "after" : "tail"); try { - await ensureAgentLoaded(msg.parentAgentId, { + await ensureUnarchivedAgentLoaded(msg.parentAgentId, { agentManager: this.agentManager, agentStorage: this.agentStorage, logger: this.sessionLogger, diff --git a/packages/server/src/server/session/agent-config/agent-config-session.test.ts b/packages/server/src/server/session/agent-config/agent-config-session.test.ts index d9ae7470c..d0e14519a 100644 --- a/packages/server/src/server/session/agent-config/agent-config-session.test.ts +++ b/packages/server/src/server/session/agent-config/agent-config-session.test.ts @@ -9,14 +9,21 @@ import type { AgentProviderNotice } from "../../agent/agent-sdk-types.js"; import type { SessionOutboundMessage } from "../../messages.js"; class FakeAgentConfigOperations implements AgentConfigOperations { + readonly loadedAgentIds: string[] = []; readonly modeCalls: Array<{ agentId: string; modeId: string }> = []; readonly modelCalls: Array<{ agentId: string; modelId: string | null }> = []; readonly featureCalls: Array<{ agentId: string; featureId: string; value: unknown }> = []; readonly thinkingCalls: Array<{ agentId: string; thinkingOptionId: string | null }> = []; modeNotice: AgentProviderNotice | null = null; thinkingNotice: AgentProviderNotice | null = null; + loadFailure: Error | null = null; failWith: Error | null = null; + async ensureLoaded(agentId: string): Promise { + this.loadedAgentIds.push(agentId); + if (this.loadFailure) throw this.loadFailure; + } + async setMode(agentId: string, modeId: string): Promise { this.modeCalls.push({ agentId, modeId }); if (this.failWith) throw this.failWith; @@ -68,6 +75,7 @@ describe("AgentConfigSession", () => { }); expect(operations.modeCalls).toEqual([{ agentId: "agent-1", modeId: "plan" }]); + expect(operations.loadedAgentIds).toEqual(["agent-1"]); expect(emitted).toEqual([ { type: "set_agent_mode_response", @@ -109,6 +117,43 @@ describe("AgentConfigSession", () => { }); }); + test("set mode: a failed load rejects without mutating the collected agent", async () => { + const { subsystem, emitted, operations } = makeSubsystem(); + operations.loadFailure = new Error("agent is archived"); + + await subsystem.handleSetAgentModeRequest({ + type: "set_agent_mode_request", + agentId: "agent-1", + modeId: "plan", + requestId: "req-1", + }); + + expect(operations.loadedAgentIds).toEqual(["agent-1"]); + expect(operations.modeCalls).toEqual([]); + expect(emitted.map((message) => message.type)).toEqual([ + "activity_log", + "set_agent_mode_response", + ]); + expect(emitted[0]).toEqual({ + type: "activity_log", + payload: { + id: expect.any(String), + timestamp: expect.any(Date), + type: "error", + content: "Failed to set agent mode: agent is archived", + }, + }); + expect(emitted[1]).toEqual({ + type: "set_agent_mode_response", + payload: { + requestId: "req-1", + agentId: "agent-1", + accepted: false, + error: "agent is archived", + }, + }); + }); + test("set model: emits an accepted response with no notice", async () => { const { subsystem, emitted, operations } = makeSubsystem(); diff --git a/packages/server/src/server/session/agent-config/agent-config-session.ts b/packages/server/src/server/session/agent-config/agent-config-session.ts index ef8630c57..1cd39856a 100644 --- a/packages/server/src/server/session/agent-config/agent-config-session.ts +++ b/packages/server/src/server/session/agent-config/agent-config-session.ts @@ -19,11 +19,12 @@ export interface AgentConfigSessionHost { /** * The per-agent config mutations this subsystem drives. The shell adapts these - * onto the live AgentManager (mode still routes through setAgentModeCommand); - * tests wire an in-memory fake. Mode and thinking yield a provider notice; model - * and feature do not. + * onto the AgentManager and loads a collected agent before mutation (mode still + * routes through setAgentModeCommand); tests wire an in-memory fake. Mode and + * thinking yield a provider notice; model and feature do not. */ export interface AgentConfigOperations { + ensureLoaded(agentId: string): Promise; setMode(agentId: string, modeId: string): Promise; setModel(agentId: string, modelId: string | null): Promise; setFeature(agentId: string, featureId: string, value: unknown): Promise; @@ -138,6 +139,7 @@ export class AgentConfigSession { this.logger.info(logFields, `session: ${logLabel}`); try { + await this.operations.ensureLoaded(agentId); const notice = await run(); this.logger.info(logFields, `session: ${logLabel} success`); emitResponse({ requestId, agentId, accepted: true, error: null, notice });