From 0e160bdab9197ab70040dc54a54a2037c01d9161 Mon Sep 17 00:00:00 2001 From: Rob Lourens Date: Mon, 24 Aug 2026 16:31:49 -0700 Subject: [PATCH] agentHost: retain recently used sessions Replace the short idle-session timeout with a soft MRU residency limit. Keep active and observed sessions pinned, release archived or least-recently-used eligible sessions non-destructively, and centralize logical subscription liveness in one service.\n\n(Written by Copilot) Co-authored-by: Copilot <223556219+Copilot@users.noreply.github.com> --- .../common/agentHostSubscriptionService.ts | 39 ++ .../platform/agentHost/common/agentService.ts | 12 +- .../agentHost/node/agentHostServices.ts | 3 + .../agentHost/node/agentHostStateManager.ts | 4 +- .../node/agentHostSubscriptionService.ts | 56 +++ .../platform/agentHost/node/agentService.ts | 339 ++++++------------ .../agentHost/node/agentServiceComposition.ts | 2 +- .../agentHost/node/agentSessionResidency.ts | 275 ++++++++++++++ .../test/node/agentHostStateManager.test.ts | 4 +- .../agentHost/test/node/agentService.test.ts | 191 ++++++++-- .../test/node/agentServiceTestUtils.ts | 4 + .../test/node/agentSessionResidency.test.ts | 267 ++++++++++++++ .../e2e/harness/agentHostE2ETestHarness.ts | 4 +- .../sessionLifecycle.integrationTest.ts | 16 +- .../test/node/protocolServerHandler.test.ts | 3 +- .../copilotMockLlm.integrationTest.ts | 29 +- 16 files changed, 945 insertions(+), 303 deletions(-) create mode 100644 src/vs/platform/agentHost/common/agentHostSubscriptionService.ts create mode 100644 src/vs/platform/agentHost/node/agentHostSubscriptionService.ts create mode 100644 src/vs/platform/agentHost/node/agentSessionResidency.ts create mode 100644 src/vs/platform/agentHost/test/node/agentSessionResidency.test.ts diff --git a/src/vs/platform/agentHost/common/agentHostSubscriptionService.ts b/src/vs/platform/agentHost/common/agentHostSubscriptionService.ts new file mode 100644 index 00000000000000..ccfe311ae02c4f --- /dev/null +++ b/src/vs/platform/agentHost/common/agentHostSubscriptionService.ts @@ -0,0 +1,39 @@ +/*--------------------------------------------------------------------------------------------- + * Copyright (c) Microsoft Corporation. All rights reserved. + * Licensed under the MIT License. See License.txt in the project root for license information. + *--------------------------------------------------------------------------------------------*/ + +import { URI } from '../../../base/common/uri.js'; +import { createDecorator } from '../../instantiation/common/instantiation.js'; +import { parseAnnotationsUri } from './annotationsUri.js'; +import { parseChangesetUri } from './changesetUri.js'; +import { parseDefaultChatUri, parseSubagentSessionUri } from './state/sessionState.js'; + +export const IAgentHostSubscriptionService = createDecorator('agentHostSubscriptionService'); + +/** + * Authoritative registry of logical protocol clients observing Agent Host state. + */ +export interface IAgentHostSubscriptionService { + readonly _serviceBrand: undefined; + + readonly subscribedResources: Iterable; + + addSubscriber(resource: URI, clientId: string): boolean; + removeSubscriber(resource: URI, clientId: string): boolean; + hasSubscribers(resource: URI): boolean; + hasSessionSubscribers(resource: URI): boolean; +} + +export function resolveAgentHostSession(resource: URI): URI { + const resourceString = resource.toString(); + const changesetSession = parseChangesetUri(resourceString)?.sessionUri; + const annotationsSession = parseAnnotationsUri(resourceString)?.sessionUri; + const chatSession = parseDefaultChatUri(resourceString); + let session = URI.parse(changesetSession ?? annotationsSession ?? chatSession ?? resourceString); + let subagent; + while ((subagent = parseSubagentSessionUri(session))) { + session = subagent.parentSession; + } + return session; +} diff --git a/src/vs/platform/agentHost/common/agentService.ts b/src/vs/platform/agentHost/common/agentService.ts index ad6340d98430d4..e35f98c52fd084 100644 --- a/src/vs/platform/agentHost/common/agentService.ts +++ b/src/vs/platform/agentHost/common/agentService.ts @@ -218,13 +218,11 @@ export const AgentHostClaudeAgentEnabledEnvVar = 'VSCODE_AGENT_HOST_CLAUDE_AGENT */ export const AgentHostCodexAgentEnabledEnvVar = 'VSCODE_AGENT_HOST_CODEX_AGENT_ENABLED'; -/** - * Overrides the grace period (in milliseconds) before an idle, fully - * unsubscribed session is released from memory. Defaults to 30_000. Primarily a - * test hook so real-SDK integration tests can force a prompt release without - * waiting the full production grace; production does not set it. - */ -export const AgentHostSessionReleaseGraceMsEnvVar = 'VSCODE_AGENT_HOST_SESSION_RELEASE_GRACE_MS'; +/** Overrides the soft cap on resident session roots. Primarily used by integration tests. */ +export const AgentHostSessionResidencyLimitEnvVar = 'VSCODE_AGENT_HOST_SESSION_RESIDENCY_LIMIT'; + +/** Overrides the retry delay after a provider temporarily vetoes session release. Primarily used by integration tests. */ +export const AgentHostSessionReleaseRetryMsEnvVar = 'VSCODE_AGENT_HOST_SESSION_RELEASE_RETRY_MS'; /** * Resolves the effective enable state for a Claude/Codex provider from the diff --git a/src/vs/platform/agentHost/node/agentHostServices.ts b/src/vs/platform/agentHost/node/agentHostServices.ts index e2f096b0f8f9b5..5b7c91b8312b59 100644 --- a/src/vs/platform/agentHost/node/agentHostServices.ts +++ b/src/vs/platform/agentHost/node/agentHostServices.ts @@ -24,6 +24,7 @@ import { IAgentHostChatContributions } from '../common/agentHostChatContribution import { IAgentHostCheckpointService } from '../common/agentHostCheckpointService.js'; import { IAgentHostGitStateService } from '../common/agentHostGitStateService.js'; import { IAgentHostReviewService } from '../common/agentHostReviewService.js'; +import { IAgentHostSubscriptionService } from '../common/agentHostSubscriptionService.js'; import { CopilotApiService, ICopilotApiService } from './shared/copilotApiService.js'; import { AgentHostFileMonitorService, IAgentHostFileMonitorService } from './agentHostFileMonitorService.js'; import { AgentHostGitService } from './agentHostGitService.js'; @@ -48,6 +49,7 @@ import { AgentHostGitStateService } from './agentHostGitStateService.js'; import { AgentHostManagedSettingsService, IAgentHostManagedSettingsService } from './agentHostManagedSettingsService.js'; import { AgentHostPromptCache, IAgentHostPromptCache } from './agentHostPromptCache.js'; import { AgentHostReviewService } from './agentHostReviewService.js'; +import { AgentHostSubscriptionService } from './agentHostSubscriptionService.js'; import { AgentHostSessionTitleSignal, IAgentHostSessionTitleSignal } from './agentHostSessionTitleSignal.js'; import { AgentHostStorageService, IAgentHostStorageService } from './agentHostStorageService.js'; import { AgentHostTerminalManager, IAgentHostTerminalManager } from './agentHostTerminalManager.js'; @@ -125,6 +127,7 @@ export function registerAgentHostCoreServices(services: AgentHostServiceCollecti registerService(services, ids, IAgentHostPromptCache, new SyncDescriptor(AgentHostPromptCache)); registerService(services, ids, IAgentHostSessionTitleSignal, new SyncDescriptor(AgentHostSessionTitleSignal)); registerService(services, ids, IAgentHostChangesetSubscriptionService, new SyncDescriptor(AgentHostChangesetSubscriptionService)); + registerService(services, ids, IAgentHostSubscriptionService, new SyncDescriptor(AgentHostSubscriptionService)); registerService(services, ids, IAgentHostChangesetOperationService, new SyncDescriptor(AgentHostChangesetOperationService)); registerService(services, ids, IAgentHostReviewService, new SyncDescriptor(AgentHostReviewService)); registerService(services, ids, IAgentHostChangesetService, new SyncDescriptor(AgentHostChangesetService)); diff --git a/src/vs/platform/agentHost/node/agentHostStateManager.ts b/src/vs/platform/agentHost/node/agentHostStateManager.ts index f02923e921f556..400286f487b80c 100644 --- a/src/vs/platform/agentHost/node/agentHostStateManager.ts +++ b/src/vs/platform/agentHost/node/agentHostStateManager.ts @@ -1186,8 +1186,8 @@ export class AgentHostStateManager extends Disposable { * be emitted as a side-effect of this call. * * Per-session changesets are intentionally NOT torn down here: this method - * is also used as an idle-eviction (LRU) hook (see - * `AgentService._maybeEvictIdleSession`) and the session list view keeps a + * is also used by `AgentSessionResidency` for residency eviction, and + * the session list view keeps a * changeset subscription open per visible row to render the diff chip. * Tearing down on eviction would clear the chip on the list while the row * is still on screen. Permanent-delete paths (`deleteSession`, diff --git a/src/vs/platform/agentHost/node/agentHostSubscriptionService.ts b/src/vs/platform/agentHost/node/agentHostSubscriptionService.ts new file mode 100644 index 00000000000000..a5d6f134215754 --- /dev/null +++ b/src/vs/platform/agentHost/node/agentHostSubscriptionService.ts @@ -0,0 +1,56 @@ +/*--------------------------------------------------------------------------------------------- + * Copyright (c) Microsoft Corporation. All rights reserved. + * Licensed under the MIT License. See License.txt in the project root for license information. + *--------------------------------------------------------------------------------------------*/ + +import { ResourceMap } from '../../../base/common/map.js'; +import { URI } from '../../../base/common/uri.js'; +import { IAgentHostSubscriptionService, resolveAgentHostSession } from '../common/agentHostSubscriptionService.js'; + +export class AgentHostSubscriptionService implements IAgentHostSubscriptionService { + declare readonly _serviceBrand: undefined; + + private readonly _subscribers = new ResourceMap>(); + + get subscribedResources(): Iterable { + return this._subscribers.keys(); + } + + addSubscriber(resource: URI, clientId: string): boolean { + let subscribers = this._subscribers.get(resource); + const firstForResource = !subscribers || subscribers.size === 0; + if (!subscribers) { + subscribers = new Set(); + this._subscribers.set(resource, subscribers); + } + subscribers.add(clientId); + return firstForResource; + } + + removeSubscriber(resource: URI, clientId: string): boolean { + const subscribers = this._subscribers.get(resource); + if (!subscribers) { + return false; + } + subscribers.delete(clientId); + if (subscribers.size > 0) { + return false; + } + this._subscribers.delete(resource); + return true; + } + + hasSubscribers(resource: URI): boolean { + return this._subscribers.has(resource); + } + + hasSessionSubscribers(resource: URI): boolean { + const sessionKey = resolveAgentHostSession(resource).toString(); + for (const subscribedResource of this._subscribers.keys()) { + if (resolveAgentHostSession(subscribedResource).toString() === sessionKey) { + return true; + } + } + return false; + } +} diff --git a/src/vs/platform/agentHost/node/agentService.ts b/src/vs/platform/agentHost/node/agentService.ts index 1138607cc8eaa7..b98380149df924 100644 --- a/src/vs/platform/agentHost/node/agentService.ts +++ b/src/vs/platform/agentHost/node/agentService.ts @@ -9,7 +9,6 @@ import { Barrier, DeferredPromise, disposableTimeout, Limiter, Promises, Resourc import { toErrorMessage } from '../../../base/common/errorMessage.js'; import { Emitter } from '../../../base/common/event.js'; import { Disposable, DisposableMap, DisposableResourceMap, DisposableStore, IDisposable, MutableDisposable } from '../../../base/common/lifecycle.js'; -import { ResourceMap } from '../../../base/common/map.js'; import { getExtensionForMimeType, getMediaMime, getMediaOrTextMime } from '../../../base/common/mime.js'; import { Schemas } from '../../../base/common/network.js'; import { ISettableObservable } from '../../../base/common/observable.js'; @@ -19,9 +18,10 @@ import { generateUuid } from '../../../base/common/uuid.js'; import { hasKey } from '../../../base/common/types.js'; import { localize } from '../../../nls.js'; import { FileChangeType, FileOperationResult, IFileChange, IFileService, toFileOperationResult, type FileChangesEvent } from '../../files/common/files.js'; +import { IInstantiationService } from '../../instantiation/common/instantiation.js'; import { ILogService } from '../../log/common/log.js'; import { AgentProvider, AgentSession, AgentSignal, IAgent, type IAgentAdoptedWorktree, IAgentChatContext, IAgentChatDataChange, IAgentChatMetadata, IAgentCreateChatOptions, IAgentCreateChatResult, IAgentCreateChatSideChatSelection, IAgentCreateChatSideChatSource, IAgentCreateSessionConfig, IAgentCreateSessionResult, IAgentDiscoveredChat, IAgentHostNetworkEndpoint, IAgentMaterializeChatEvent, IAgentModelInfo, IAgentResolveSessionConfigParams, IAgentChatAdoptionResult, type AgentChatAdoptionReason, IAgentSessionConfigCompletionsParams, IAgentSessionMetadata, IAgentSpawnChatEvent, AuthenticateParams, AuthenticateResult, IMcpNotification, SubagentChatSignal, subagentChatTitle } from '../common/agent.js'; -import { AgentHostSessionReleaseGraceMsEnvVar, type AgentHostDebugLogsArtifactKind, type IAgentHostDebugLogsArtifact, type IAgentHostDebugLogsChunk, IAgentHostManagedSettingsDiagnostics, IAgentHostNetworkDiagnosticsInfo, IAgentHostNetworkFetchResult, IAgentService } from '../common/agentService.js'; +import { type AgentHostDebugLogsArtifactKind, type IAgentHostDebugLogsArtifact, type IAgentHostDebugLogsChunk, IAgentHostManagedSettingsDiagnostics, IAgentHostNetworkDiagnosticsInfo, IAgentHostNetworkFetchResult, IAgentService } from '../common/agentService.js'; import { ISessionDataService, SESSION_ATTACHMENTS_DIRNAME } from '../common/sessionDataService.js'; import { IAgentEditAttributionService, ICancelEditAttributionFlushParams, ICommitEditAttributionFlushParams, IEditAttributionFlushResult, IPrepareEditAttributionFlushParams, IPreparedEditAttributionFlush, parseEditAttributionResource } from '../common/fileEditAttribution.js'; import { SessionConfigKey } from '../common/sessionConfigKeys.js'; @@ -53,8 +53,10 @@ import { AgentHostDebugLogsCollector, type IAgentHostDebugLogsEnvironment } from import { IAgentHostDatabase } from './agentHostDatabase.js'; import { AgentSessionRegistry, IRegisteredSession, IStoredRegisteredSession } from './agentSessionRegistry.js'; import { IAgentHostGitService } from '../common/agentHostGitService.js'; +import { IAgentHostSubscriptionService, resolveAgentHostSession } from '../common/agentHostSubscriptionService.js'; import { AgentSideEffects, type IAgentSideEffectsOptions } from './agentSideEffects.js'; import { AgentHostLocalTurns } from './agentHostLocalTurns.js'; +import { AgentSessionResidency } from './agentSessionResidency.js'; import { AgentServerToolHost } from './shared/agentServerToolHost.js'; import { type IChatContextSnapshot, type IRenameTitleResult, type ISessionCreationDefaults, type ISessionServerToolAccessor, validateRenameTitle } from './shared/sessionServerTools.js'; import { AGENT_HOST_TITLE_SOURCE_AGENT, customChatTitleMetadataKey, customChatTitleSourceMetadataKey, persistSessionMetadata, persistSessionMetadataValues, SESSION_ARTIFACTS_KEY, SESSION_CUSTOM_TITLE_KEY, SESSION_CUSTOM_TITLE_SOURCE_KEY } from './shared/persistSessionMetadata.js'; @@ -185,18 +187,6 @@ const RESOURCE_WATCH_GRACE_MS = 30_000; /** Bound on how long {@link AgentService.subscribe} waits for a pending subagent chat to register before giving up. */ const SUBAGENT_CHAT_PENDING_TIMEOUT_MS = 15_000; -/** - * Grace period before an idle session is released from memory via - * {@link AgentService._maybeEvictIdleSession}. This lets a quick reconnect - * reuse the live SDK session instead of forcing an immediate release/resume - * cycle. Overridable via {@link AgentHostSessionReleaseGraceMsEnvVar} in tests. - */ -const SESSION_RELEASE_GRACE_MS = (() => { - const raw = process.env[AgentHostSessionReleaseGraceMsEnvVar]; - const parsed = raw !== undefined ? parseInt(raw, 10) : NaN; - return Number.isFinite(parsed) && parsed >= 0 ? parsed : 30_000; -})(); - /** * Session-database metadata key for the orchestrator-owned catalog of * additional peer chats. When absent, the session predates this persistence @@ -364,6 +354,8 @@ export interface IAgentServiceOptions { readonly storageResource?: URI; readonly orchestratorDatabase?: IAgentHostDatabase; readonly debugLogsEnvironment?: IAgentHostDebugLogsEnvironment; + readonly sessionResidencyLimit?: number; + readonly sessionReleaseRetryMs?: number; } export interface IAgentServiceCallbacks { @@ -542,13 +534,11 @@ export class AgentService extends Disposable implements IAgentService { * client IDs. Populated by {@link subscribe} (or {@link addSubscriber} * for handshake fast-paths) and drained by {@link unsubscribe}. When a * resource's set becomes empty, the resource is dropped from the map and - * {@link _maybeEvictIdleSession} is invoked to release any cached state - * for it. + * session residency is reconciled against the MRU cap. */ - private readonly _resourceSubscribers = new ResourceMap>(); - private readonly _releaseSessionInFlight = new Map>(); private readonly _restoreSessionInFlight = new Map>(); private readonly _restoreSubagentInFlight = new Map>(); + private readonly _sessionResidency: AgentSessionResidency; /** * Persisted-annotation reads in flight, keyed by session URI. Annotations @@ -571,15 +561,6 @@ export class AgentService extends Disposable implements IAgentService { */ private readonly _pendingSessionGc = this._register(new DisposableResourceMap()); - /** - * Pending {@link _maybeEvictIdleSession} timers, keyed by session URI. A - * timer is armed when an idle session (with turns) loses its last subscriber - * — see {@link unsubscribe}. Cleared when any client subscribes again - * ({@link addSubscriber}) or the timer fires. Deferring the release avoids - * churning the provider SDK session on rapid disconnect/reconnect cycles. - */ - private readonly _pendingSessionRelease = this._register(new DisposableResourceMap()); - /** * Active resource watches keyed by the channel URI string * (`ahp-resource-watch:/`). @@ -607,14 +588,17 @@ export class AgentService extends Disposable implements IAgentService { constructor( core: IAgentServiceCore, collaborators: IAgentServiceCollaborators, + options: IAgentServiceOptions, @ILogService private readonly _logService: ILogService, @IFileService private readonly _fileService: IFileService, @ISessionDataService private readonly _sessionDataService: ISessionDataService, @IAgentHostGitService private readonly _gitService: IAgentHostGitService, @ITelemetryService private readonly _telemetryService: ITelemetryService, @IAgentHostChatContributions private readonly _chatContributions: IAgentHostChatContributions, + @IAgentHostSubscriptionService private readonly _subscriptions: IAgentHostSubscriptionService, @INetworkDiagnosticsService private readonly _networkDiagnostics: INetworkDiagnosticsService, @IAgentEditAttributionService private readonly _editAttributionService: IAgentEditAttributionService, + @IInstantiationService instantiationService: IInstantiationService, ) { super(); this._authService = core.authenticationService; @@ -639,6 +623,29 @@ export class AgentService extends Disposable implements IAgentService { this._sideEffects = collaborators.sideEffects; this._sessionCoordination = collaborators.sessionCoordination; this._serverToolHost = collaborators.serverToolHost; + this._sessionResidency = this._register(instantiationService.createInstance( + AgentSessionResidency, + this._stateManager, + { + isReleaseBlocked: session => this._restoreSessionInFlight.has(session.toString()), + whenSessionDataIdle: session => this._whenSessionDataIdle(session), + getSessionChats: session => this._getSessionChatsInTeardownOrder(session), + createRelease: session => { + const provider = this._findProviderForSession(session); + return provider ? { + canRelease: chats => this._canReleaseSession(provider, session, chats), + release: chats => this._releaseSession(provider, session, chats), + } : undefined; + }, + evictSessionState: (session, chats) => this._evictSessionState(session, session.toString(), session.toString(), chats.map(chat => chat.toString())), + }, + { + limit: options.sessionResidencyLimit, + releaseRetryMs: options.sessionReleaseRetryMs, + holdsSession: session => this._agentMergeController.holdsSession(session), + onDidReleaseHold: this._agentMergeController.onDidReleaseHold, + }, + )); core.callbackBinder.bind({ canEvictChangeset: changeset => this._canEvictChangeset(changeset), startAgentMergeTurn: (session, turnId, prompt) => this._startAgentMergePrompt(session, turnId, prompt), @@ -661,6 +668,7 @@ export class AgentService extends Disposable implements IAgentService { this._register(this._stateManager.onDidEmitEnvelope(e => { if (e.action.type === ActionType.SessionIsArchivedChanged && e.action.isArchived && !isAhpChatChannel(e.channel)) { this._clearAgentMergeIndex(URI.parse(e.channel)); + void this._sessionResidency.reconcile(); } })); this._register(this._stateManager.onDidEmitNotification(e => this._onDidNotification.fire(e))); @@ -675,13 +683,6 @@ export class AgentService extends Disposable implements IAgentService { })); updateAgentHostTelemetryLevelFromConfig(this._telemetryService, this._stateManager.rootState.config?.values); this._register(this._stateManager.onDidChangeSessionConfig(({ session, previous, current }) => this._syncAgentMergeIndex(URI.parse(session), previous, current))); - this._register(this._agentMergeController.onDidReleaseHold(session => { - const resource = URI.parse(session); - if (!this._hasSessionSubscribers(resource) && this._stateManager.getSessionState(session)) { - this._scheduleSessionRelease(resource); - } - })); - let externalSessionsMode = this._getExternalSessionsMode(); this._lastMigrateLegacyEnabled = this._isMigrateLegacyEnabled(); let agentMergeEnabled = this._isAgentMergeEnabled(); @@ -1401,14 +1402,6 @@ export class AgentService extends Disposable implements IAgentService { } this._logService.info(`[AgentService] Restoring Agent-Merge-enabled session for monitoring: ${sessionStr}`); await this.restoreSession(session); - // `restoreSession` cancels any pending release and arms none, so - // a session that never takes the hold (e.g. a stale index entry - // whose config says disabled) would otherwise stay resident with - // nothing left to release it. While the hold does apply, this - // timer simply finds the session held and stands down. - if (!this._hasSessionSubscribers(session) && this._stateManager.getSessionState(sessionStr)) { - this._scheduleSessionRelease(session); - } } catch (err) { this._logService.warn(`[AgentService] Failed to restore Agent-Merge-enabled session ${sessionStr}`, err); } @@ -2521,7 +2514,7 @@ export class AgentService extends Disposable implements IAgentService { } if (config?.session) { this._cancelPendingSessionGc(config.session); - this._cancelPendingSessionRelease(config.session); + this._sessionResidency.touch(config.session); } // Capability gate: only a provider that advertises @@ -2594,7 +2587,7 @@ export class AgentService extends Disposable implements IAgentService { // timer would still fire and dispose the just-revived session // before the follow-up `subscribe` arrives. this._cancelPendingSessionGc(session); - this._cancelPendingSessionRelease(session); + this._sessionResidency.touch(session); this._logService.trace(`[AgentService] createSession: provider=${provider.id} model=${config?.model?.id ?? '(default)'}`); this._sessionToProvider.set(session.toString(), provider.id); @@ -2742,6 +2735,8 @@ export class AgentService extends Disposable implements IAgentService { void this._gitStateService.refreshSessionGitState(session.toString(), workingDirectory); } + this._sessionResidency.touch(session); + await this._sessionResidency.reconcile(); return session; } @@ -2852,6 +2847,8 @@ export class AgentService extends Disposable implements IAgentService { ...(providerData !== undefined ? { providerData } : {}), ...(peerChatOrigin !== undefined ? { origin: peerChatOrigin } : {}), }); + this._sessionResidency.touch(session); + void this._sessionResidency.reconcile(); // If the agent exposes this chat as its own SDK session, mark that // backing so it stays out of the top-level session list. `_markChatBacking` @@ -3039,13 +3036,29 @@ export class AgentService extends Disposable implements IAgentService { private _getSessionChatsInTeardownOrder(session: URI): URI[] { const state = this._stateManager.getSessionState(session.toString()); - const defaultChat = state?.defaultChat ?? buildDefaultChatUri(session.toString()); + return this._orderSessionChatsForTeardown(session, state?.chats.map(chat => chat.resource) ?? []); + } + + private async _getSessionChatsForDisposal(provider: IAgent, session: URI): Promise { + const state = this._stateManager.getSessionState(session.toString()); + if (state) { + return this._getSessionChatsInTeardownOrder(session); + } + const persisted = await this._readPersistedPeerChatCatalog(session); + const peerChats = persisted?.map(chat => chat.uri) + ?? (await provider.listLegacyChatBackings?.(session))?.map(chat => chat.uri.toString()) + ?? []; + return this._orderSessionChatsForTeardown(session, peerChats); + } + + private _orderSessionChatsForTeardown(session: URI, chats: readonly string[]): URI[] { + const defaultChat = buildDefaultChatUri(session.toString()); const result: URI[] = []; const seen = new Set(); - for (const summary of state?.chats ?? []) { - if (summary.resource !== defaultChat && !seen.has(summary.resource)) { - seen.add(summary.resource); - result.push(URI.parse(summary.resource)); + for (const chat of chats) { + if (chat !== defaultChat && !seen.has(chat)) { + seen.add(chat); + result.push(URI.parse(chat)); } } if (!seen.has(defaultChat)) { @@ -3061,7 +3074,7 @@ export class AgentService extends Disposable implements IAgentService { private async _disposeSession(provider: IAgent, session: URI): Promise { await this._defaultChatBackingWrites.get(session.toString())?.catch(() => { }); let firstError: unknown; - for (const chat of this._getSessionChatsInTeardownOrder(session)) { + for (const chat of await this._getSessionChatsForDisposal(provider, session)) { try { await provider.chats.disposeChat(chat, this._chatContext(session, chat)); } catch (err) { @@ -3659,7 +3672,12 @@ export class AgentService extends Disposable implements IAgentService { async disposeSession(session: URI): Promise { this._logService.trace(`[AgentService] disposeSession: ${session.toString()}`); + await this._sessionResidency.runDisposal(session, () => this._doDisposeSession(session)); + } + + private async _doDisposeSession(session: URI): Promise { const sessionKey = session.toString(); + this._cancelPendingSessionGc(session); const isEphemeral = this._stateManager.isEphemeralSession(sessionKey); this._stateManager.invalidateSessionChatResolutions(session.toString()); const sessionChats = this._stateManager.getSessionState(session.toString())?.chats ?? []; @@ -3739,7 +3757,7 @@ export class AgentService extends Disposable implements IAgentService { this._logService.trace(`[AgentService] subscribe: ${resource.toString()}`); const resourceStr = resource.toString(); try { - await this._releaseSessionInFlight.get(this._sessionReleaseKey(resource)); + await this._sessionResidency.waitForRelease(resource); // Register after an in-flight release settles so a successful release // can evict cached state and this subscribe reconstructs it. The // handshake fast path calls addSubscriber directly and therefore pins @@ -3821,6 +3839,8 @@ export class AgentService extends Disposable implements IAgentService { if (!snapshot) { throw new Error(`Cannot subscribe to unknown resource: ${resourceStr}`); } + this._sessionResidency.touch(resource); + void this._sessionResidency.reconcile(); // Ensure git state has been computed for this session. When the snapshot // already existed (e.g. seeded by list query, or restored earlier), the @@ -3847,18 +3867,6 @@ export class AgentService extends Disposable implements IAgentService { } } - private _sessionReleaseKey(resource: URI): string { - const resourceString = resource.toString(); - const changesetSession = parseChangesetUri(resourceString)?.sessionUri; - const chatSession = parseDefaultChatUri(resourceString); - let session = URI.parse(changesetSession ?? chatSession ?? resourceString); - let subagent; - while ((subagent = parseSubagentSessionUri(session))) { - session = subagent.parentSession; - } - return session.toString(); - } - /** Waits for an armed subagent chat to register (or its wait to time out); returns `undefined` if not armed or never registered. */ private async _awaitPendingSubagentChat(subagentChatUri: string): Promise { const pending = this._pendingSubagentChats.get(subagentChatUri); @@ -3870,37 +3878,24 @@ export class AgentService extends Disposable implements IAgentService { } addSubscriber(resource: URI, clientId: string): void { - let set = this._resourceSubscribers.get(resource); - const wasUnsubscribed = !set || set.size === 0; - if (!set) { - set = new Set(); - this._resourceSubscribers.set(resource, set); - } - set.add(clientId); // A new subscriber means the session is being observed again; cancel // any pending GC or idle-release armed while it had no subscribers. this._cancelPendingSessionGc(resource); this._cancelPendingEphemeralSessionGc(resource); - this._cancelPendingSessionRelease(resource); // 0→1 transition — covers both the full subscribe path AND the // handshake fast-path used by `ProtocolServerHandler` when state is // already cached. The coordinator decides whether the URI is one // it cares about (e.g. uncommitted changeset → trigger refresh). - if (wasUnsubscribed) { + if (this._subscriptions.addSubscriber(resource, clientId)) { this._changesetCoordinator.onFirstSubscriber(resource); } + this._sessionResidency.touch(resource); } unsubscribe(resource: URI, clientId: string): void { - const set = this._resourceSubscribers.get(resource); - if (!set) { - return; - } - set.delete(clientId); - if (set.size > 0) { + if (!this._subscriptions.removeSubscriber(resource, clientId)) { return; } - this._resourceSubscribers.delete(resource); this._changesetCoordinator.onLastSubscriber(resource); this._stateManager.onChangesetLivenessChanged(); if (this._maybeScheduleEphemeralSessionGc(resource)) { @@ -3908,21 +3903,15 @@ export class AgentService extends Disposable implements IAgentService { } // An empty session whose last subscriber dropped is a candidate for // full GC (provider session, worktree, on-disk state). Sessions with - // at least one turn fall through to {@link _maybeEvictIdleSession}, - // which only drops the in-memory cache and lets the session be - // restored from disk later. Skipping eviction here for empty + // at least one turn participate in residency reconciliation, which only + // drops the in-memory cache and lets the session be restored from disk + // later. Skipping eviction here for empty // sessions ensures their state stays observable so a re-subscribe // can re-arm GC. if (this._maybeScheduleSessionGc(resource)) { return; } - // Defer the idle-session release behind a grace window rather than - // releasing synchronously. A client that reconnects (or re-subscribes) - // within the window cancels this via {@link _cancelPendingSessionRelease} - // and keeps the live provider SDK session, avoiding a disconnect/resume - // churn cycle that races concurrent session operations on the shared - // provider runtime. A zero grace releases on the next tick. - this._scheduleSessionRelease(resource); + void this._sessionResidency.reconcile(); } /** @@ -3930,12 +3919,12 @@ export class AgentService extends Disposable implements IAgentService { * chat subscriptions are gone, regardless of whether it has completed turns. */ private _maybeScheduleEphemeralSessionGc(resource: URI): boolean { - const session = this._sessionReleaseResource(resource); + const session = resolveAgentHostSession(resource); const sessionKey = session.toString(); if (!this._stateManager.isEphemeralSession(sessionKey)) { return false; } - if (this._hasSessionSubscribers(session)) { + if (this._subscriptions.hasSessionSubscribers(session)) { return true; } this._pendingSessionGc.set(session, disposableTimeout(() => { @@ -3947,30 +3936,12 @@ export class AgentService extends Disposable implements IAgentService { return true; } - private _cancelPendingSessionRelease(resource: URI): void { - this._pendingSessionRelease.deleteAndDispose(this._sessionReleaseResource(resource)); - } - - private _scheduleSessionRelease(resource: URI): void { - const session = this._sessionReleaseResource(resource); - this._pendingSessionRelease.set(session, disposableTimeout(() => { - this._pendingSessionRelease.deleteAndDispose(session); - void this._maybeEvictIdleSession(session).catch(err => { - this._logService.error(err, `[AgentService] Failed to evict idle session ${session.toString()}`); - }); - }, SESSION_RELEASE_GRACE_MS)); - } - - private _sessionReleaseResource(resource: URI): URI { - return URI.parse(this._sessionReleaseKey(resource)); - } - /** * If `resource` names a session that no client is still subscribed to and * that has produced no turns (and has no active turn), schedule a delayed * {@link _runSessionGc} to fully tear it down — provider session, worktree, * persisted state and all. Sessions with at least one turn are left to the - * existing {@link _maybeEvictIdleSession} path which only drops cached + * residency path which only drops cached * state and lets the session be restored from disk later. * * GC is restricted to sessions that are still unused drafts. A session that @@ -3987,12 +3958,11 @@ export class AgentService extends Disposable implements IAgentService { * so callers can skip alternative cleanup paths. */ private _maybeScheduleSessionGc(resource: URI): boolean { - // Subagent URIs are backed by the parent session; the parent's GC is - // scheduled when its own subscriber count reaches zero. - if (parseSubagentSessionUri(resource)) { - return false; + const session = resolveAgentHostSession(resource); + if (this._subscriptions.hasSessionSubscribers(session)) { + return true; } - const key = resource.toString(); + const key = session.toString(); const state = this._stateManager.getSessionState(key); if (!state) { return false; @@ -4005,12 +3975,12 @@ export class AgentService extends Disposable implements IAgentService { return false; } // Never tear down a session Agent Merge is holding. - if (this._agentMergeController.holdsSession(this._sessionReleaseKey(resource))) { + if (this._agentMergeController.holdsSession(key)) { return false; } - this._pendingSessionGc.set(resource, disposableTimeout(() => { - this._pendingSessionGc.deleteAndDispose(resource); - this._runSessionGc(resource).catch(err => { + this._pendingSessionGc.set(session, disposableTimeout(() => { + this._pendingSessionGc.deleteAndDispose(session); + this._runSessionGc(session).catch(err => { this._logService.error(err, `[AgentService] GC failed for ${key}`); }); }, SESSION_GC_GRACE_MS)); @@ -4018,18 +3988,18 @@ export class AgentService extends Disposable implements IAgentService { } private _cancelPendingSessionGc(resource: URI): void { - this._pendingSessionGc.deleteAndDispose(resource); + this._pendingSessionGc.deleteAndDispose(resolveAgentHostSession(resource)); } private _cancelPendingEphemeralSessionGc(resource: URI): void { - const session = this._sessionReleaseResource(resource); + const session = resolveAgentHostSession(resource); if (this._stateManager.isEphemeralSession(session.toString())) { this._pendingSessionGc.deleteAndDispose(session); } } private async _runEphemeralSessionGc(session: URI): Promise { - if (this._hasSessionSubscribers(session)) { + if (this._subscriptions.hasSessionSubscribers(session)) { return; } this._logService.info(`[AgentService] GC: disposing unsubscribed ephemeral session ${session.toString()}`); @@ -4041,12 +4011,12 @@ export class AgentService extends Disposable implements IAgentService { * subscriber while empty. Re-checks the invariants (still no subscribers, * still empty, still an unused draft) before tearing the session down via * {@link disposeSession}. The cached state may already have been evicted by - * {@link _maybeEvictIdleSession}; in that case we still proceed because + * residency reconciliation; in that case we still proceed because * "evicted + no resubscribe" implies no client is observing the session. */ private async _runSessionGc(resource: URI): Promise { const key = resource.toString(); - if (this._resourceSubscribers.has(resource)) { + if (this._subscriptions.hasSessionSubscribers(resource)) { return; } const state = this._stateManager.getSessionState(key); @@ -4064,107 +4034,6 @@ export class AgentService extends Disposable implements IAgentService { await this.disposeSession(resource); } - /** - * If `resource` names an idle session with no remaining subscribers, drop its - * cached state and release its SDK chats. Subagent URIs evict the parent - * session entry because the parent owns the materialized turn tree. Durable - * data stays intact; the next subscribe restores the session on demand. - */ - private async _maybeEvictIdleSession(resource: URI): Promise { - const key = resource.toString(); - const evictionTarget = this._sessionReleaseResource(resource); - const evictionTargetKey = evictionTarget.toString(); - if (this._hasSessionSubscribers(evictionTarget)) { - return; - } - // A restore/resume racing this unsubscribe means a client is about to - // observe the session again; releasing now would tear down state that - // the in-flight rehydrate is populating. - if (this._restoreSessionInFlight.has(evictionTargetKey)) { - return; - } - const targetState = this._stateManager.getSessionState(evictionTargetKey); - if (!targetState) { - return; - } - if (this._stateManager.hasActiveTurn(evictionTargetKey)) { - this._scheduleSessionRelease(evictionTarget); - return; - } - // Agent Merge keeps monitoring with no client subscriber, so releasing - // would silently stop it until someone reopened the session. - if (this._agentMergeController.holdsSession(evictionTargetKey)) { - this._logService.trace(`[AgentService] Skipping idle eviction for a session held by Agent Merge: ${evictionTargetKey}`); - return; - } - if (this._releaseSessionInFlight.has(evictionTargetKey)) { - return; - } - const chats = this._getSessionChatsInTeardownOrder(evictionTarget); - await this._whenSessionDataIdle(evictionTarget); - if (this._hasSessionSubscribers(evictionTarget) || this._restoreSessionInFlight.has(evictionTargetKey) || this._releaseSessionInFlight.has(evictionTargetKey)) { - return; - } - const settledState = this._stateManager.getSessionState(evictionTargetKey); - if (!settledState) { - return; - } - if (this._stateManager.hasActiveTurn(evictionTargetKey)) { - this._scheduleSessionRelease(evictionTarget); - return; - } - if (this._agentMergeController.holdsSession(evictionTargetKey)) { - return; - } - const provider = this._findProviderForSession(evictionTarget); - if (!provider) { - return; - } - const trackedRelease = (async () => { - try { - if (!await this._canReleaseSession(provider, evictionTarget, chats)) { - if (!this._hasSessionSubscribers(evictionTarget)) { - this._scheduleSessionRelease(evictionTarget); - } - return; - } - const currentState = this._stateManager.getSessionState(evictionTargetKey); - if (this._hasSessionSubscribers(evictionTarget)) { - return; - } - if (this._restoreSessionInFlight.has(evictionTargetKey) || this._stateManager.hasActiveTurn(evictionTargetKey)) { - this._scheduleSessionRelease(evictionTarget); - return; - } - if (currentState) { - this._evictSessionState(evictionTarget, evictionTargetKey, key, currentState.chats.map(chat => chat.resource)); - } - await this._releaseSession(provider, evictionTarget, chats); - } catch (err) { - this._logService.error(err, `[AgentService] Failed to release idle session ${evictionTargetKey}`); - if (!this._hasSessionSubscribers(evictionTarget)) { - this._scheduleSessionRelease(evictionTarget); - } - } - })(); - this._releaseSessionInFlight.set(evictionTargetKey, trackedRelease); - void trackedRelease.then(() => { - if (this._releaseSessionInFlight.get(evictionTargetKey) === trackedRelease) { - this._releaseSessionInFlight.delete(evictionTargetKey); - } - }); - } - - private _hasSessionSubscribers(session: URI): boolean { - const sessionKey = this._sessionReleaseKey(session); - for (const subscribedUri of this._resourceSubscribers.keys()) { - if (this._sessionReleaseKey(subscribedUri) === sessionKey) { - return true; - } - } - return false; - } - private _evictSessionState(evictionTarget: URI, evictionTargetKey: string, triggerKey: string, chats: readonly string[]): void { this._logService.info(`[AgentService] Evicting idle session: ${evictionTargetKey} (triggered by unsubscribe of ${triggerKey})`); const subagentPrefix = buildSubagentSessionUriPrefix(evictionTarget); @@ -4181,7 +4050,7 @@ export class AgentService extends Disposable implements IAgentService { const changesetUri = URI.parse(changeset); // A direct changeset subscriber is rendering this expanded URI. Keep // the state alive so future envelopes still target an existing object. - if (this._resourceSubscribers.has(changesetUri)) { + if (this._subscriptions.hasSubscribers(changesetUri)) { return false; } const parsed = parseChangesetUri(changeset); @@ -4192,12 +4061,12 @@ export class AgentService extends Disposable implements IAgentService { const sessionUri = URI.parse(parsed.sessionUri); // A parent-session subscriber can still receive catalogue count updates // from this changeset, so keep the backing state while the session is observed. - if (this._resourceSubscribers.has(sessionUri)) { + if (this._subscriptions.hasSubscribers(sessionUri)) { return false; } // Subagent views are backed by the parent session tree; treat any // subscribed descendant as a parent-session pin for cache eviction. - for (const subscribedUri of this._resourceSubscribers.keys()) { + for (const subscribedUri of this._subscriptions.subscribedResources) { if (this._isSubagentDescendantOf(subscribedUri, sessionUri)) { return false; } @@ -4657,8 +4526,8 @@ export class AgentService extends Disposable implements IAgentService { async restoreSession(session: URI): Promise { const sessionStr = session.toString(); this._cancelPendingSessionGc(session); - this._cancelPendingSessionRelease(session); - await this._releaseSessionInFlight.get(sessionStr); + this._sessionResidency.touch(session); + await this._sessionResidency.waitForRelease(session); const inFlight = this._restoreSessionInFlight.get(sessionStr); if (inFlight) { @@ -4667,6 +4536,8 @@ export class AgentService extends Disposable implements IAgentService { } if (this._stateManager.getSessionState(sessionStr)) { + this._sessionResidency.touch(session); + await this._sessionResidency.reconcile(); return; } @@ -4675,6 +4546,8 @@ export class AgentService extends Disposable implements IAgentService { this._restoreSessionInFlight.set(sessionStr, restore); try { await restore; + this._sessionResidency.touch(session); + await this._sessionResidency.reconcile(); this._logService.trace(`[AgentService] restoreSession done: ${sessionStr}`); } finally { if (this._restoreSessionInFlight.get(sessionStr) === restore) { @@ -6395,7 +6268,7 @@ export class AgentService extends Disposable implements IAgentService { if (!this._gitService) { throw new ProtocolError(AhpErrorCodes.NotFound, `git service unavailable for: ${fields.repoRelativePath}`); } - const owningSession = this._sessionReleaseResource(URI.parse(fields.sessionUri)); + const owningSession = resolveAgentHostSession(URI.parse(fields.sessionUri)); const wasRestored = !!this._stateManager.getSessionState(owningSession.toString()); try { if (!wasRestored) { @@ -6415,8 +6288,8 @@ export class AgentService extends Disposable implements IAgentService { contentType: 'text/plain', }; } finally { - if (!wasRestored && this._stateManager.getSessionState(owningSession.toString()) && !this._hasSessionSubscribers(owningSession)) { - this._scheduleSessionRelease(owningSession); + if (!wasRestored && this._stateManager.getSessionState(owningSession.toString()) && !this._subscriptions.hasSessionSubscribers(owningSession)) { + void this._sessionResidency.reconcile(); } } } diff --git a/src/vs/platform/agentHost/node/agentServiceComposition.ts b/src/vs/platform/agentHost/node/agentServiceComposition.ts index 4211e22d3b2ad6..4e772ebc37fd61 100644 --- a/src/vs/platform/agentHost/node/agentServiceComposition.ts +++ b/src/vs/platform/agentHost/node/agentServiceComposition.ts @@ -166,7 +166,7 @@ export function createAgentServiceComposition( sessionCoordination, serverToolHost, }; - agentService = instantiationService.createInstance(AgentService, core, collaborators); + agentService = instantiationService.createInstance(AgentService, core, collaborators, options); for (const disposable of additionalDisposables) { owned.add(disposable); } diff --git a/src/vs/platform/agentHost/node/agentSessionResidency.ts b/src/vs/platform/agentHost/node/agentSessionResidency.ts new file mode 100644 index 00000000000000..aab6aac9b5c620 --- /dev/null +++ b/src/vs/platform/agentHost/node/agentSessionResidency.ts @@ -0,0 +1,275 @@ +/*--------------------------------------------------------------------------------------------- + * Copyright (c) Microsoft Corporation. All rights reserved. + * Licensed under the MIT License. See License.txt in the project root for license information. + *--------------------------------------------------------------------------------------------*/ + +import { disposableTimeout, Sequencer } from '../../../base/common/async.js'; +import type { Event } from '../../../base/common/event.js'; +import { Disposable, DisposableResourceMap, type IDisposable } from '../../../base/common/lifecycle.js'; +import { LinkedMap, Touch } from '../../../base/common/map.js'; +import { URI } from '../../../base/common/uri.js'; +import { ILogService } from '../../log/common/log.js'; +import { AgentHostSessionReleaseRetryMsEnvVar, AgentHostSessionResidencyLimitEnvVar } from '../common/agentService.js'; +import { IAgentHostSubscriptionService, resolveAgentHostSession } from '../common/agentHostSubscriptionService.js'; +import { isSessionStatusArchived, parseSubagentSessionUri } from '../common/state/sessionState.js'; +import { AgentHostStateManager } from './agentHostStateManager.js'; + +const DEFAULT_SESSION_RESIDENCY_LIMIT = 10; +const DEFAULT_SESSION_RELEASE_RETRY_MS = 30_000; + +function readNonNegativeIntegerEnv(name: string, defaultValue: number): number { + const raw = process.env[name]; + const parsed = raw !== undefined ? parseInt(raw, 10) : NaN; + return Number.isFinite(parsed) && parsed >= 0 ? parsed : defaultValue; +} + +export interface IAgentSessionRelease { + canRelease(chats: readonly URI[]): Promise; + release(chats: readonly URI[]): Promise; +} + +export interface IAgentSessionReleaseDelegate { + isReleaseBlocked(session: URI): boolean; + whenSessionDataIdle(session: URI): Promise; + getSessionChats(session: URI): readonly URI[]; + createRelease(session: URI): IAgentSessionRelease | undefined; + evictSessionState(session: URI, chats: readonly URI[]): void; +} + +export interface IAgentSessionResidencyOptions { + readonly limit?: number; + readonly releaseRetryMs?: number; + readonly holdsSession: (session: string) => boolean; + readonly onDidReleaseHold: Event; +} + +export class AgentSessionResidency extends Disposable { + private readonly _releaseInFlight = new Map>(); + private readonly _sessionsBeingDisposed = new Set(); + private readonly _recency = new LinkedMap(); + private readonly _reconciler = new Sequencer(); + private readonly _releaseRetries = this._register(new DisposableResourceMap()); + private readonly _limit: number; + private readonly _releaseRetryMs: number; + + constructor( + private readonly _stateManager: AgentHostStateManager, + private readonly _delegate: IAgentSessionReleaseDelegate, + private readonly _options: IAgentSessionResidencyOptions, + @IAgentHostSubscriptionService private readonly _subscriptions: IAgentHostSubscriptionService, + @ILogService private readonly _logService: ILogService, + ) { + super(); + this._limit = _options.limit ?? readNonNegativeIntegerEnv(AgentHostSessionResidencyLimitEnvVar, DEFAULT_SESSION_RESIDENCY_LIMIT); + this._releaseRetryMs = _options.releaseRetryMs ?? readNonNegativeIntegerEnv(AgentHostSessionReleaseRetryMsEnvVar, DEFAULT_SESSION_RELEASE_RETRY_MS); + this._register(this._stateManager.onDidChangeSessionActiveTurn(({ session, active }) => { + if (active) { + this.touch(URI.parse(session)); + } + void this.reconcile(); + })); + this._register(this._stateManager.onDidRemoveSession(session => { + const resource = URI.parse(session); + if (!parseSubagentSessionUri(resource)) { + this._recency.delete(session); + void this.reconcile(); + } + })); + this._register(this._options.onDidReleaseHold(() => void this.reconcile())); + } + + touch(resource: URI): void { + const session = resolveAgentHostSession(resource); + this._releaseRetries.deleteAndDispose(session); + const sessionKey = session.toString(); + if (!this._stateManager.getSessionState(sessionKey) + || this._stateManager.isEphemeralSession(sessionKey) + || this._stateManager.isUnusedDraft(sessionKey) === true) { + return; + } + this._recency.set(sessionKey, session, Touch.AsNew); + } + + reconcile(): Promise { + return this._reconciler.queue(() => this._doReconcile()); + } + + async runDisposal(resource: URI, task: () => Promise): Promise { + const session = resolveAgentHostSession(resource); + const sessionKey = session.toString(); + this._sessionsBeingDisposed.add(sessionKey); + this._recency.delete(sessionKey); + this._releaseRetries.deleteAndDispose(session); + try { + await this._releaseInFlight.get(sessionKey); + this._releaseRetries.deleteAndDispose(session); + try { + return await task(); + } catch (error) { + this.touch(session); + throw error; + } + } finally { + this._sessionsBeingDisposed.delete(sessionKey); + void this.reconcile(); + } + } + + async waitForRelease(resource: URI): Promise { + await this._releaseInFlight.get(resolveAgentHostSession(resource).toString()); + } + + private _residentCount(): number { + let count = 0; + for (const [sessionKey] of this._recency) { + if (this._stateManager.getSessionState(sessionKey) + && !this._stateManager.isEphemeralSession(sessionKey) + && this._stateManager.isUnusedDraft(sessionKey) !== true) { + count++; + } + } + return count; + } + + private _isReleaseRequired(sessionKey: string, expectedRecency: URI | undefined): boolean { + if (isSessionStatusArchived(this._stateManager.getSessionState(sessionKey)?.status)) { + return true; + } + return expectedRecency !== undefined + && this._recency.get(sessionKey) === expectedRecency + && this._residentCount() > this._limit; + } + + private async _doReconcile(): Promise { + const stale: string[] = []; + const residents: URI[] = []; + for (const [sessionKey, session] of this._recency) { + if (!this._stateManager.getSessionState(sessionKey) + || this._stateManager.isEphemeralSession(sessionKey) + || this._stateManager.isUnusedDraft(sessionKey) === true) { + stale.push(sessionKey); + } else { + residents.push(session); + } + } + for (const sessionKey of stale) { + this._recency.delete(sessionKey); + } + + const archived = new Map(); + for (const sessionKey of this._stateManager.getSessionUris()) { + const session = URI.parse(sessionKey); + if (!parseSubagentSessionUri(session) + && isSessionStatusArchived(this._stateManager.getSessionState(sessionKey)?.status)) { + archived.set(sessionKey, session); + } + } + + let residentCount = residents.length; + const processed = new Set(); + for (const session of [...archived.values(), ...residents]) { + const sessionKey = session.toString(); + if (processed.has(sessionKey)) { + continue; + } + processed.add(sessionKey); + if (!archived.has(sessionKey) && residentCount <= this._limit) { + break; + } + const wasResident = this._recency.has(sessionKey); + if (await this._tryRelease(session, this._recency.get(sessionKey)) && wasResident) { + residentCount--; + } + } + } + + private async _tryRelease(resource: URI, expectedRecency: URI | undefined): Promise { + const session = resolveAgentHostSession(resource); + const sessionKey = session.toString(); + if (!this._isReleaseRequired(sessionKey, expectedRecency) + || this._sessionsBeingDisposed.has(sessionKey) + || this._releaseRetries.has(session) + || this._subscriptions.hasSessionSubscribers(session) + || this._delegate.isReleaseBlocked(session) + || this._stateManager.hasActiveTurn(sessionKey) + || this._options.holdsSession(sessionKey)) { + return false; + } + const releaseInFlight = this._releaseInFlight.get(sessionKey); + if (releaseInFlight) { + return releaseInFlight; + } + const trackedRelease = this._release(session, expectedRecency); + this._releaseInFlight.set(sessionKey, trackedRelease); + try { + return await trackedRelease; + } finally { + if (this._releaseInFlight.get(sessionKey) === trackedRelease) { + this._releaseInFlight.delete(sessionKey); + } + } + } + + private async _release(session: URI, expectedRecency: URI | undefined): Promise { + const sessionKey = session.toString(); + try { + await this._delegate.whenSessionDataIdle(session); + if (!this._canContinueRelease(sessionKey, expectedRecency)) { + return false; + } + + const release = this._delegate.createRelease(session); + if (!release) { + return false; + } + let chats = this._delegate.getSessionChats(session); + while (true) { + if (!await release.canRelease(chats)) { + this._scheduleRetryIfNeeded(session, expectedRecency); + return false; + } + const currentChats = this._delegate.getSessionChats(session); + if (currentChats.length === chats.length && currentChats.every((chat, index) => chat.toString() === chats[index].toString())) { + break; + } + chats = currentChats; + } + + if (!this._stateManager.getSessionState(sessionKey) || !this._canContinueRelease(sessionKey, expectedRecency)) { + return false; + } + await release.release(chats); + if (!this._canContinueRelease(sessionKey, expectedRecency)) { + return false; + } + this._delegate.evictSessionState(session, chats); + return true; + } catch (error) { + this._logService.error(error, `[AgentSessionResidency] Failed to release session ${sessionKey}`); + this._scheduleRetryIfNeeded(session, expectedRecency); + return false; + } + } + + private _canContinueRelease(sessionKey: string, expectedRecency: URI | undefined): boolean { + return this._isReleaseRequired(sessionKey, expectedRecency) + && !this._sessionsBeingDisposed.has(sessionKey) + && !this._subscriptions.hasSessionSubscribers(URI.parse(sessionKey)) + && !this._delegate.isReleaseBlocked(URI.parse(sessionKey)) + && !this._stateManager.hasActiveTurn(sessionKey) + && !this._options.holdsSession(sessionKey); + } + + private _scheduleRetryIfNeeded(session: URI, expectedRecency: URI | undefined): void { + const sessionKey = session.toString(); + if (!this._isReleaseRequired(sessionKey, expectedRecency) + || this._sessionsBeingDisposed.has(sessionKey) + || this._subscriptions.hasSessionSubscribers(session)) { + return; + } + this._releaseRetries.set(session, disposableTimeout(() => { + this._releaseRetries.deleteAndDispose(session); + void this.reconcile(); + }, this._releaseRetryMs)); + } +} diff --git a/src/vs/platform/agentHost/test/node/agentHostStateManager.test.ts b/src/vs/platform/agentHost/test/node/agentHostStateManager.test.ts index f250d5e427218a..9bfd5a5e895ccd 100644 --- a/src/vs/platform/agentHost/test/node/agentHostStateManager.test.ts +++ b/src/vs/platform/agentHost/test/node/agentHostStateManager.test.ts @@ -881,7 +881,7 @@ suite('AgentHostStateManager', () => { }); test('removeSession flushes pending status=Idle notification before eviction', () => { - // Regression: when _maybeEvictIdleSession calls removeSession within the + // Regression: when residency eviction calls removeSession within the // 100 ms scheduler window after a turn completes, the client must still // receive a SessionSummaryChanged with status=Idle so the spinner clears. // @@ -961,7 +961,7 @@ suite('AgentHostStateManager', () => { }); test('removeSession does NOT dispose per-session changesets (LRU eviction must not clear list-view chip)', () => { - // Regression: _maybeEvictIdleSession calls removeSession to drop an + // Regression: residency eviction calls removeSession to drop an // idle session from the in-memory cache. The Agents Window list view // keeps a per-row changeset subscription open to render the diff // chip, so cascading disposeSessionChangesets here would emit a diff --git a/src/vs/platform/agentHost/test/node/agentService.test.ts b/src/vs/platform/agentHost/test/node/agentService.test.ts index a2cb4864e707f1..6013e940098c0b 100644 --- a/src/vs/platform/agentHost/test/node/agentService.test.ts +++ b/src/vs/platform/agentHost/test/node/agentService.test.ts @@ -1429,7 +1429,7 @@ suite('AgentService (node dispatcher)', () => { }); }); - test('git-blob temporarily restores the owning session for nested subagents', () => { + test('git-blob restores the owning session for nested subagents and retains it in the MRU', () => { return runWithFakedTimers({ useFakeTimers: true }, async () => { const repoA = URI.file('/workspace/repoA'); const showBlobCalls: Array<{ workingDirectory: string; ref: string; repoRelativePath: string }> = []; @@ -1459,7 +1459,7 @@ suite('AgentService (node dispatcher)', () => { showBlobCalls: [{ workingDirectory: repoA.toString(), ref: 'baseSha', repoRelativePath: 'src/app.ts' }], data: 'blob:src/app.ts', retainedForSubscriber: true, - releasedAfterUnsubscribe: true, + releasedAfterUnsubscribe: false, }); }); }); @@ -11481,7 +11481,59 @@ suite('AgentService (node dispatcher)', () => { }); }); - suite('subscriber refcount eviction', () => { + suite('session residency eviction', () => { + + function createResidencyTestService(limit: number, releaseRetryMs = 30_000, registerProvider = true, agent = new MockAgent('copilot'), sessionDataService = nullSessionDataService): { readonly service: AgentService; readonly agent: MockAgent } { + const testService = disposables.add(createTestAgentService( + new NullLogService(), + fileService, + sessionDataService, + { _serviceBrand: undefined } as IProductService, + createNoopGitService(), + undefined, + undefined, + undefined, + undefined, + globalThis.fetch, + [], + undefined, + undefined, + undefined, + limit, + releaseRetryMs, + )); + disposables.add(toDisposable(() => agent.dispose())); + if (registerProvider) { + testService.registerProvider(agent); + } + return { service: testService, agent }; + } + + async function createUsedSession(testService: AgentService, agent: MockAgent, complete = true): Promise { + const session = await testService.createSession({ provider: agent.id }); + const chat = buildDefaultChatUri(session); + getStateManager(testService).dispatchServerAction(chat, { type: ActionType.ChatTurnStarted, turnId: `turn-${AgentSession.id(session)}`, startedAt: '2025-01-01T00:00:00.000Z', message: { text: 'hello', origin: { kind: MessageKind.User } } }); + if (complete) { + getStateManager(testService).dispatchServerAction(chat, { type: ActionType.ChatTurnComplete, turnId: `turn-${AgentSession.id(session)}`, duration: 1000 }); + } + return session; + } + + async function waitForResidency(predicate: () => boolean, message: string): Promise { + for (let attempt = 0; attempt < 100; attempt++) { + if (predicate()) { + return; + } + await timeout(0); + } + assert.fail(message); + } + + setup(() => { + const zeroCapacity = createResidencyTestService(0, 30_000, false); + service = zeroCapacity.service; + copilotAgent = zeroCapacity.agent; + }); class DelayedReleaseMockAgent extends MockAgent { readonly release = new DeferredPromise(); @@ -11532,6 +11584,56 @@ suite('AgentService (node dispatcher)', () => { })); } + test('deleting an evicted session disposes every persisted peer chat', async () => { + class MultiChatAgent extends MockAgent { + override async createChat(_session: URI, _chat: URI): Promise { } + } + const agent = new MultiChatAgent('copilot'); + const residency = createResidencyTestService(1, 30_000, true, agent); + const first = await createUsedSession(residency.service, residency.agent); + const peer = URI.parse(buildChatUri(first, 'peer-1')); + await residency.service.createChat(first, peer); + const second = await createUsedSession(residency.service, residency.agent); + await waitForResidency(() => getStateManager(residency.service).getSessionState(first.toString()) === undefined, 'first session was not evicted'); + + await residency.service.disposeSession(first); + + assert.deepStrictEqual({ + residentSession: getStateManager(residency.service).getSessionState(second.toString()) !== undefined, + disposedChats: residency.agent.chatContexts.filter(call => call.boundary === 'disposeChat').map(call => call.chat.toString()), + }, { + residentSession: true, + disposedChats: [peer.toString(), buildDefaultChatUri(first)], + }); + }); + + test('serializes deletion behind an in-flight release preflight', async () => { + const whenIdleStarted = new DeferredPromise(); + const whenIdle = new DeferredPromise(); + class DelayedIdleDatabase extends TestSessionDatabase { + override async whenIdle(): Promise { + whenIdleStarted.complete(); + await whenIdle.p; + } + } + const residency = createResidencyTestService(10, 30_000, true, new MockAgent('copilot'), createSessionDataService(new DelayedIdleDatabase())); + const session = await createUsedSession(residency.service, residency.agent); + getStateManager(residency.service).dispatchServerAction(session.toString(), { type: ActionType.SessionIsArchivedChanged, isArchived: true }); + await whenIdleStarted.p; + + const deletion = residency.service.disposeSession(session); + whenIdle.complete(); + await deletion; + + assert.deepStrictEqual({ + releases: residency.agent.releaseSessionCalls.map(call => call.toString()), + disposals: residency.agent.disposeSessionCalls.map(call => call.toString()), + }, { + releases: [], + disposals: [session.toString()], + }); + }); + test('an empty session created in this lifetime stays observable until GC fires', async () => { service.registerProvider(copilotAgent); const sessionResource = await service.createSession({ provider: 'copilot' }); @@ -11602,7 +11704,23 @@ suite('AgentService (node dispatcher)', () => { await whenIdle.p; } } - const localService = disposables.add(createTestAgentService(new NullLogService(), fileService, createSessionDataService(new DelayedIdleDatabase()), { _serviceBrand: undefined } as IProductService, createNoopGitService())); + const localService = disposables.add(createTestAgentService( + new NullLogService(), + fileService, + createSessionDataService(new DelayedIdleDatabase()), + { _serviceBrand: undefined } as IProductService, + createNoopGitService(), + undefined, + undefined, + undefined, + undefined, + globalThis.fetch, + [], + undefined, + undefined, + undefined, + 0, + )); const agent = new MockAgent('copilot'); disposables.add(toDisposable(() => agent.dispose())); localService.registerProvider(agent); @@ -11615,7 +11733,6 @@ suite('AgentService (node dispatcher)', () => { localService.addSubscriber(sessionResource, 'client-1'); localService.unsubscribe(sessionResource, 'client-1'); - await new Promise(resolve => setTimeout(resolve, 30_000)); await whenIdleStarted.p; localService.dispatchAction( peerChat.toString(), @@ -11668,7 +11785,7 @@ suite('AgentService (node dispatcher)', () => { }); }); - test('chat subscription cancels the root release retry and gets a fresh grace period', () => { + test('chat subscription cancels the root release retry until the chat unsubscribes', () => { return runWithFakedTimers({ useFakeTimers: true }, async () => { const agent = new DeferringReleaseMockAgent('copilot'); service.registerProvider(agent); @@ -11686,13 +11803,10 @@ suite('AgentService (node dispatcher)', () => { assert.strictEqual(agent.releaseAttempts, 1); service.addSubscriber(chatResource, 'client-chat'); - await new Promise(resolve => setTimeout(resolve, 10_000)); + await new Promise(resolve => setTimeout(resolve, 30_000)); + assert.strictEqual(agent.releaseAttempts, 1, 'the cancelled retry must not fire while a chat is subscribed'); service.unsubscribe(chatResource, 'client-chat'); - await new Promise(resolve => setTimeout(resolve, 20_000)); - assert.strictEqual(agent.releaseAttempts, 1, 'the cancelled root retry must not fire at its original deadline'); - await new Promise(resolve => setTimeout(resolve, 9_999)); - assert.strictEqual(agent.releaseAttempts, 1, 'chat disconnect should receive a fresh release grace'); - await new Promise(resolve => setTimeout(resolve, 1)); + await new Promise(resolve => setTimeout(resolve, 0)); assert.deepStrictEqual({ releaseAttempts: agent.releaseAttempts, @@ -11704,7 +11818,7 @@ suite('AgentService (node dispatcher)', () => { }); }); - test('overlapping root release timers preserve the original in-flight release', () => { + test('overlapping residency reconciliations preserve the original in-flight release', () => { return runWithFakedTimers({ useFakeTimers: true }, async () => { const agent = new DelayedReleaseMockAgent('copilot'); service.registerProvider(agent); @@ -11725,7 +11839,7 @@ suite('AgentService (node dispatcher)', () => { service.addSubscriber(chatResource, 'client-chat'); service.unsubscribe(chatResource, 'client-chat'); await new Promise(resolve => setTimeout(resolve, 30_000)); - assert.deepStrictEqual(agent.events, ['release:start'], 'second timer must not start another provider release'); + assert.deepStrictEqual(agent.events, ['release:start'], 'second reconciliation must not start another provider release'); await agent.release.complete(); await Promise.resolve(); @@ -11748,11 +11862,11 @@ suite('AgentService (node dispatcher)', () => { service.addSubscriber(sessionResource, 'client-1'); service.unsubscribe(sessionResource, 'client-1'); - // Release is deferred behind the grace window — still cached until it elapses. - assert.ok(getStateManager(service).getSessionState(sessionResource.toString()), 'session stays cached during the release grace'); + // Reconciliation is asynchronous, so state remains cached until release preflight runs. + assert.ok(getStateManager(service).getSessionState(sessionResource.toString()), 'session stays cached until release preflight'); await new Promise(resolve => setTimeout(resolve, 30_000)); - assert.strictEqual(getStateManager(service).getSessionState(sessionResource.toString()), undefined, 'restored idle session should be evicted after the grace'); + assert.strictEqual(getStateManager(service).getSessionState(sessionResource.toString()), undefined, 'restored idle session should be evicted after reconciliation'); assert.deepStrictEqual( copilotAgent.releaseSessionCalls.map(u => u.toString()), [sessionResource.toString()], @@ -11798,7 +11912,7 @@ suite('AgentService (node dispatcher)', () => { }); }); - test('re-subscribing within the grace cancels the release', () => { + test('re-subscribing during release preflight cancels the release', () => { return runWithFakedTimers({ useFakeTimers: true }, async () => { service.registerProvider(copilotAgent); const { session } = await createAgentSession(copilotAgent); @@ -11813,12 +11927,12 @@ suite('AgentService (node dispatcher)', () => { service.addSubscriber(sessionResource, 'client-1'); service.unsubscribe(sessionResource, 'client-1'); - // Reconnect within the grace window. + // Reconnect before asynchronous release preflight completes. service.addSubscriber(sessionResource, 'client-2'); await new Promise(resolve => setTimeout(resolve, 30_000)); - assert.ok(getStateManager(service).getSessionState(sessionResource.toString()), 'session must stay cached when re-subscribed within the grace'); - assert.strictEqual(copilotAgent.releaseSessionCalls.length, 0, 'chat release must not fire when the grace was cancelled'); + assert.ok(getStateManager(service).getSessionState(sessionResource.toString()), 'session must stay cached when re-subscribed during release preflight'); + assert.strictEqual(copilotAgent.releaseSessionCalls.length, 0, 'chat release must not fire after preflight was cancelled'); }); }); @@ -12185,6 +12299,30 @@ suite('AgentService (node dispatcher)', () => { suite('empty-session GC', () => { + test('a default-chat subscriber pins an empty session after the root unsubscribes', () => { + return runWithFakedTimers({ useFakeTimers: true }, async () => { + service.registerProvider(copilotAgent); + const session = await service.createSession({ provider: 'copilot' }); + const chat = URI.parse(buildDefaultChatUri(session)); + service.addSubscriber(session, 'session-client'); + service.addSubscriber(chat, 'chat-client'); + + service.unsubscribe(session, 'session-client'); + await new Promise(resolve => setTimeout(resolve, 30_000)); + const residentForChat = getStateManager(service).getSessionState(session.toString()) !== undefined; + service.unsubscribe(chat, 'chat-client'); + await new Promise(resolve => setTimeout(resolve, 30_000)); + + assert.deepStrictEqual({ + residentForChat, + disposals: copilotAgent.disposeSessionCalls.map(call => call.toString()), + }, { + residentForChat: true, + disposals: [session.toString()], + }); + }); + }); + test('an empty unsubscribed session is disposed after the grace period', () => { return runWithFakedTimers({ useFakeTimers: true }, async () => { service.registerProvider(copilotAgent); @@ -12309,9 +12447,8 @@ suite('AgentService (node dispatcher)', () => { disposed: copilotAgent.disposeSessionCalls.map(u => u.toString()), released: copilotAgent.releaseSessionCalls.map(u => u.toString()), }, { - // Nothing destroyed, but the non-destructive idle release still happens. disposed: [], - released: [sessionResource.toString()], + released: [], }); }); }); @@ -13220,7 +13357,7 @@ suite('AgentService (node dispatcher)', () => { return { localService, localAgent, sessionResource }; } - test('an Agent-Merge-enabled session stays resident after its last subscriber drops, and is released once disabled', () => { + test('an Agent-Merge-enabled session stays resident after its last subscriber drops and becomes MRU-eligible once disabled', () => { return runWithFakedTimers({ useFakeTimers: true }, async () => { const orchestratorDb = new TestAgentHostOrchestratorDatabase(); const { localService, sessionResource } = await createEnabledSession(new TestSessionDatabase(), orchestratorDb); @@ -13240,7 +13377,7 @@ suite('AgentService (node dispatcher)', () => { indexedAfterDisable: await orchestratorDb.listAgentMergeEnabledSessions(), }, { residentWhileEnabled: true, - residentAfterDisable: false, + residentAfterDisable: true, indexedAfterDisable: [], }); }); @@ -13253,7 +13390,7 @@ suite('AgentService (node dispatcher)', () => { assert.deepStrictEqual(await orchestratorDb.listAgentMergeEnabledSessions(), [sessionResource.toString()]); }); - test('a persisted Agent-Merge-enabled session begins monitoring on a fresh host, and is released once disabled', () => { + test('a persisted Agent-Merge-enabled session begins monitoring on a fresh host and becomes MRU-eligible once disabled', () => { return runWithFakedTimers({ useFakeTimers: true }, async () => { const orchestratorDb = new TestAgentHostOrchestratorDatabase(); const sessionDb = new TestSessionDatabase(); @@ -13283,7 +13420,7 @@ suite('AgentService (node dispatcher)', () => { residentAfterDisable: getStateManager(restarted).getSessionState(sessionStr) !== undefined, }, { resumed: { materialized: true, enabled: true, indexed: [sessionStr] }, - residentAfterDisable: false, + residentAfterDisable: true, }); }); }); diff --git a/src/vs/platform/agentHost/test/node/agentServiceTestUtils.ts b/src/vs/platform/agentHost/test/node/agentServiceTestUtils.ts index e8899d1718e10d..7cb2fdedf74ba3 100644 --- a/src/vs/platform/agentHost/test/node/agentServiceTestUtils.ts +++ b/src/vs/platform/agentHost/test/node/agentServiceTestUtils.ts @@ -79,6 +79,8 @@ export function createTestAgentService( hostLaunchKind = AgentHostLaunchKind.Unknown, storageResource?: URI, orchestratorDatabase?: IAgentHostDatabase, + sessionResidencyLimit?: number, + sessionReleaseRetryMs?: number, ): AgentService { const effectiveFileMonitorService = fileMonitorService ?? new AgentHostFileMonitorService(fileService, logService); const clientConnectionService = new AgentHostClientConnectionService(); @@ -102,6 +104,8 @@ export function createTestAgentService( hostLaunchKind, storageResource, orchestratorDatabase, + sessionResidencyLimit, + sessionReleaseRetryMs, }; const foundationDisposables = new DisposableStore(); const foundation = createAgentServiceFoundation({ diff --git a/src/vs/platform/agentHost/test/node/agentSessionResidency.test.ts b/src/vs/platform/agentHost/test/node/agentSessionResidency.test.ts new file mode 100644 index 00000000000000..c9e22b62b7d97d --- /dev/null +++ b/src/vs/platform/agentHost/test/node/agentSessionResidency.test.ts @@ -0,0 +1,267 @@ +/*--------------------------------------------------------------------------------------------- + * Copyright (c) Microsoft Corporation. All rights reserved. + * Licensed under the MIT License. See License.txt in the project root for license information. + *--------------------------------------------------------------------------------------------*/ + +import assert from 'assert'; +import { DeferredPromise, timeout } from '../../../../base/common/async.js'; +import { Emitter } from '../../../../base/common/event.js'; +import { DisposableStore } from '../../../../base/common/lifecycle.js'; +import { URI } from '../../../../base/common/uri.js'; +import { ensureNoDisposablesAreLeakedInTestSuite } from '../../../../base/test/common/utils.js'; +import { ILogService, NullLogService } from '../../../log/common/log.js'; +import { InstantiationService } from '../../../instantiation/common/instantiationService.js'; +import { ServiceCollection } from '../../../instantiation/common/serviceCollection.js'; +import { AgentSession } from '../../common/agent.js'; +import { IAgentHostSubscriptionService } from '../../common/agentHostSubscriptionService.js'; +import { buildAnnotationsUri } from '../../common/annotationsUri.js'; +import { ActionType } from '../../common/state/sessionActions.js'; +import { MessageKind, SessionStatus, buildDefaultChatUri, type SessionSummary } from '../../common/state/sessionState.js'; +import { AgentHostStateManager } from '../../node/agentHostStateManager.js'; +import { AgentHostSubscriptionService } from '../../node/agentHostSubscriptionService.js'; +import { AgentSessionResidency, type IAgentSessionReleaseDelegate } from '../../node/agentSessionResidency.js'; + +suite('AgentSessionResidency', () => { + const disposables = new DisposableStore(); + const logService = new NullLogService(); + let stateManager: AgentHostStateManager; + let releaseHold: Emitter; + let released: string[]; + let evicted: string[]; + let residency: AgentSessionResidency; + let delegate: IAgentSessionReleaseDelegate; + let subscriptions: AgentHostSubscriptionService; + + setup(() => { + stateManager = disposables.add(new AgentHostStateManager(logService)); + releaseHold = disposables.add(new Emitter()); + released = []; + evicted = []; + subscriptions = new AgentHostSubscriptionService(); + delegate = { + isReleaseBlocked: () => false, + whenSessionDataIdle: async () => { }, + getSessionChats: session => [URI.parse(buildDefaultChatUri(session))], + createRelease: session => ({ + canRelease: async () => true, + release: async () => { released.push(session.toString()); }, + }), + evictSessionState: session => { + evicted.push(session.toString()); + stateManager.removeSession(session.toString()); + }, + }; + residency = createResidency(10); + }); + + teardown(() => disposables.clear()); + ensureNoDisposablesAreLeakedInTestSuite(); + + function createResidency(limit: number, releaseRetryMs = 30_000): AgentSessionResidency { + const instantiationService = disposables.add(new InstantiationService(new ServiceCollection( + [ILogService, logService], + [IAgentHostSubscriptionService, subscriptions], + ))); + return disposables.add(instantiationService.createInstance( + AgentSessionResidency, + stateManager, + delegate, + { + limit, + releaseRetryMs, + holdsSession: () => false, + onDidReleaseHold: releaseHold.event, + }, + )); + } + + function createUsedSession(id: string, complete = true): URI { + const session = AgentSession.uri('copilot', id); + const summary: SessionSummary = { + resource: session.toString(), + provider: 'copilot', + title: id, + status: SessionStatus.Idle, + createdAt: '2025-01-01T00:00:00.000Z', + modifiedAt: '2025-01-01T00:00:00.000Z', + }; + stateManager.createSession(summary); + const chat = buildDefaultChatUri(session); + stateManager.dispatchServerAction(chat, { type: ActionType.ChatTurnStarted, turnId: `turn-${id}`, startedAt: '2025-01-01T00:00:00.000Z', message: { text: id, origin: { kind: MessageKind.User } } }); + if (complete) { + stateManager.dispatchServerAction(chat, { type: ActionType.ChatTurnComplete, turnId: `turn-${id}`, duration: 1 }); + } + residency.touch(session); + return session; + } + + async function waitFor(predicate: () => boolean, message: string): Promise { + for (let attempt = 0; attempt < 100; attempt++) { + if (predicate()) { + return; + } + await timeout(0); + } + assert.fail(message); + } + + test('keeps the most-recently-used sessions within the soft limit', async () => { + residency.dispose(); + residency = createResidency(3); + const first = createUsedSession('first'); + const second = createUsedSession('second'); + const third = createUsedSession('third'); + residency.touch(first); + const fourth = createUsedSession('fourth'); + + await residency.reconcile(); + + assert.deepStrictEqual({ + resident: [first, second, third, fourth].map(session => stateManager.getSessionState(session.toString()) !== undefined), + released, + evicted, + }, { + resident: [true, false, true, true], + released: [second.toString()], + evicted: [second.toString()], + }); + }); + + test('allows pinned sessions to exceed the limit and trims when one becomes idle', async () => { + residency.dispose(); + residency = createResidency(2); + const first = createUsedSession('first', false); + const second = createUsedSession('second', false); + const third = createUsedSession('third', false); + + await residency.reconcile(); + const residentWhileRunning = [first, second, third].map(session => stateManager.getSessionState(session.toString()) !== undefined); + stateManager.dispatchServerAction(buildDefaultChatUri(first), { type: ActionType.ChatTurnComplete, turnId: 'turn-first', duration: 1 }); + await waitFor(() => stateManager.getSessionState(first.toString()) === undefined, 'completed excess session was not released'); + + assert.deepStrictEqual({ + residentWhileRunning, + residentAfterCompletion: [first, second, third].map(session => stateManager.getSessionState(session.toString()) !== undefined), + }, { + residentWhileRunning: [true, true, true], + residentAfterCompletion: [false, true, true], + }); + }); + + test('releases archived sessions below the limit', async () => { + const session = createUsedSession('archived'); + stateManager.dispatchServerAction(session.toString(), { type: ActionType.SessionIsArchivedChanged, isArchived: true }); + + await residency.reconcile(); + + assert.deepStrictEqual({ + resident: stateManager.getSessionState(session.toString()) !== undefined, + released, + }, { + resident: false, + released: [session.toString()], + }); + }); + + test('normalizes annotations subscriptions to their owning session', async () => { + residency.dispose(); + residency = createResidency(0); + const session = createUsedSession('annotations'); + const annotations = URI.parse(buildAnnotationsUri(session.toString())); + subscriptions.addSubscriber(annotations, 'client'); + + await residency.reconcile(); + const residentForSubscriber = stateManager.getSessionState(session.toString()) !== undefined; + const removedLastSubscriber = subscriptions.removeSubscriber(annotations, 'client'); + const hasSubscribersAfterRemoval = subscriptions.hasSessionSubscribers(session); + + assert.deepStrictEqual({ + removedLastSubscriber, + hasSubscribersAfterRemoval, + residentForSubscriber, + }, { + removedLastSubscriber: true, + hasSubscribersAfterRemoval: false, + residentForSubscriber: true, + }); + }); + + test('keeps state after a failed release and retries', async () => { + residency.dispose(); + let attempts = 0; + delegate.createRelease = session => ({ + canRelease: async () => true, + release: async () => { + attempts++; + if (attempts === 1) { + throw new Error('transient'); + } + released.push(session.toString()); + }, + }); + residency = createResidency(10, 10); + const session = createUsedSession('retry'); + stateManager.dispatchServerAction(session.toString(), { type: ActionType.SessionIsArchivedChanged, isArchived: true }); + + await residency.reconcile(); + const residentAfterFailure = stateManager.getSessionState(session.toString()) !== undefined; + await timeout(15); + await waitFor(() => stateManager.getSessionState(session.toString()) === undefined, 'release retry did not complete'); + + assert.deepStrictEqual({ + attempts, + residentAfterFailure, + }, { + attempts: 2, + residentAfterFailure: true, + }); + }); + + test('restores MRU tracking when disposal fails', async () => { + residency.dispose(); + residency = createResidency(0); + const session = createUsedSession('dispose-failure'); + + await assert.rejects( + () => residency.runDisposal(session, async () => { throw new Error('transient disposal failure'); }), + /transient disposal failure/, + ); + await waitFor(() => stateManager.getSessionState(session.toString()) === undefined, 'failed disposal did not restore residency tracking'); + + assert.deepStrictEqual(released, [session.toString()]); + }); + + test('revalidates recency after asynchronous provider preflight', async () => { + residency.dispose(); + const canRelease = new DeferredPromise(); + let createReleaseCalls = 0; + delegate.createRelease = session => { + createReleaseCalls++; + return { + canRelease: async () => { + await canRelease.p; + return true; + }, + release: async () => { released.push(session.toString()); }, + }; + }; + residency = createResidency(1); + const first = createUsedSession('first'); + const second = createUsedSession('second'); + const reconcile = residency.reconcile(); + await timeout(0); + + residency.touch(first); + canRelease.complete(); + await reconcile; + await residency.reconcile(); + + assert.deepStrictEqual({ + createReleaseCalls, + resident: [first, second].map(session => stateManager.getSessionState(session.toString()) !== undefined), + }, { + createReleaseCalls: 2, + resident: [true, false], + }); + }); +}); diff --git a/src/vs/platform/agentHost/test/node/e2e/harness/agentHostE2ETestHarness.ts b/src/vs/platform/agentHost/test/node/e2e/harness/agentHostE2ETestHarness.ts index 622501ce6bd762..09072b227dccd1 100644 --- a/src/vs/platform/agentHost/test/node/e2e/harness/agentHostE2ETestHarness.ts +++ b/src/vs/platform/agentHost/test/node/e2e/harness/agentHostE2ETestHarness.ts @@ -30,7 +30,7 @@ import { type ChatErrorAction, type ChatToolCallCompleteAction, type ChatToolCallStartAction, } from '../../../../common/state/sessionActions.js'; import { CopilotCliConfigKey } from '../../../../common/copilotCliConfig.js'; -import { AgentHostSessionReleaseGraceMsEnvVar } from '../../../../common/agentService.js'; +import { AgentHostSessionResidencyLimitEnvVar } from '../../../../common/agentService.js'; import { CapiReplayMode, type ICapiReplayResponse } from './capiReplayProxy.js'; import { fetchSessionWithChat, getActionEnvelope, getAgentHostE2ETestTimeout, isActionNotification, IServerHandle, stopServer, TestProtocolClient, @@ -879,7 +879,7 @@ export class AgentHostE2EServerLease { codexHomeDir, homeDir: dataDir, userDataDir: join(dataDir, 'user-data'), - env: { [AgentHostSessionReleaseGraceMsEnvVar]: '0' }, + env: { [AgentHostSessionResidencyLimitEnvVar]: '0' }, }; // Server reuse is a replay-only optimization: recording writes one fixture // per proxy and so needs a fresh proxy (hence a fresh server) per test. diff --git a/src/vs/platform/agentHost/test/node/protocol/sessionLifecycle.integrationTest.ts b/src/vs/platform/agentHost/test/node/protocol/sessionLifecycle.integrationTest.ts index 601dec1dd2383d..794874c7168018 100644 --- a/src/vs/platform/agentHost/test/node/protocol/sessionLifecycle.integrationTest.ts +++ b/src/vs/platform/agentHost/test/node/protocol/sessionLifecycle.integrationTest.ts @@ -13,7 +13,7 @@ import type { SessionAddedParams, SessionRemovedParams } from '../../../common/s import { PROTOCOL_VERSION } from '../../../common/state/protocol/version/registry.js'; import type { ListSessionsResult } from '../../../common/state/sessionProtocol.js'; import { buildDefaultChatUri, ResponsePartKind, ROOT_STATE_URI, SessionStatus, type MarkdownResponsePart, type ISessionWithDefaultChat, type ToolCallResponsePart } from '../../../common/state/sessionState.js'; -import { AgentHostSessionReleaseGraceMsEnvVar } from '../../../common/agentService.js'; +import { AgentHostSessionResidencyLimitEnvVar } from '../../../common/agentService.js'; import { AgentHostExternalSessionsMode, AgentHostShowExternalSessionsConfigKey } from '../../../common/agentHostSchema.js'; import { PRE_EXISTING_SESSION_URI } from '../mockAgent.js'; import { @@ -41,10 +41,7 @@ suite('Protocol WebSocket — Session Lifecycle', function () { return secondaryClient; } - // Short idle-release grace so the release/restore test exercises a real - // release promptly. Safe on the shared server because the mock agent's - // releaseSession is cheap (no real SDK disconnect). - const RELEASE_GRACE_MS = 200; + const RESIDENCY_SETTLE_MS = 500; suiteSetup(async function () { this.timeout(getAgentHostE2ETestTimeout(15_000, 60_000)); @@ -54,7 +51,7 @@ suite('Protocol WebSocket — Session Lifecycle', function () { env: { HOME: userDataDir, USERPROFILE: userDataDir, - [AgentHostSessionReleaseGraceMsEnvVar]: String(RELEASE_GRACE_MS), + [AgentHostSessionResidencyLimitEnvVar]: '0', }, }); }); @@ -192,12 +189,11 @@ suite('Protocol WebSocket — Session Lifecycle', function () { const before = await fetchSessionWithChat(client, preExistingUri); assert.ok(before.turns.length >= 1, 'session should restore turns on first subscribe'); - // Drop every subscriber; after the release grace elapses the server - // evicts the idle session (dropping cached state and releasing the - // provider's SDK resources). + // Drop every subscriber; the zero-capacity test server evicts the idle + // session and releases the provider's SDK resources. client.notify('unsubscribe', { channel: chatUri }); client.notify('unsubscribe', { channel: preExistingUri }); - await timeout(RELEASE_GRACE_MS + 500); + await timeout(RESIDENCY_SETTLE_MS); // Re-subscribing rehydrates the session from the preserved durable data // — the turns must match the pre-eviction view. diff --git a/src/vs/platform/agentHost/test/node/protocolServerHandler.test.ts b/src/vs/platform/agentHost/test/node/protocolServerHandler.test.ts index adc343a46ba54e..536546d6833635 100644 --- a/src/vs/platform/agentHost/test/node/protocolServerHandler.test.ts +++ b/src/vs/platform/agentHost/test/node/protocolServerHandler.test.ts @@ -2344,8 +2344,7 @@ suite('ProtocolServerHandler', () => { transport1.simulateClose(); // Simulate the AgentService evicting the idle session while the client - // was disconnected (this is what `_maybeEvictIdleSession` does in the - // real service). + // was disconnected (this is what residency eviction does in the real service). stateManager.removeSession(sessionUri); assert.strictEqual(stateManager.getSnapshot(sessionUri), undefined, 'precondition: state evicted'); diff --git a/src/vs/platform/agentHost/test/node/providerIntegration/copilotMockLlm.integrationTest.ts b/src/vs/platform/agentHost/test/node/providerIntegration/copilotMockLlm.integrationTest.ts index 54d3907f41245a..b0e71fe70ac615 100644 --- a/src/vs/platform/agentHost/test/node/providerIntegration/copilotMockLlm.integrationTest.ts +++ b/src/vs/platform/agentHost/test/node/providerIntegration/copilotMockLlm.integrationTest.ts @@ -18,7 +18,7 @@ import { URI } from '../../../../../base/common/uri.js'; import { ActionType, type ChatToolCallCompleteAction, type ChatToolCallReadyAction } from '../../../common/state/sessionActions.js'; import { buildDefaultChatUri, ResponsePartKind, SessionStatus, type ISessionWithDefaultChat } from '../../../common/state/sessionState.js'; import { ToolCallConfirmationReason } from '../../../common/state/protocol/channels-chat/state.js'; -import { AgentHostSessionReleaseGraceMsEnvVar } from '../../../common/agentService.js'; +import { AgentHostSessionReleaseRetryMsEnvVar, AgentHostSessionResidencyLimitEnvVar } from '../../../common/agentService.js'; import { createProviderSession, dispatchTurn, type IAgentHostProviderTestConfig } from '../providerIntegrationTestHelpers.js'; import { fetchSessionWithChat, getActionEnvelope, isActionNotification, IServerHandle, startRealServer, stopServer, TestProtocolClient } from '../serverIntegrationTestHelpers.js'; @@ -109,17 +109,14 @@ suite('Agent Host Provider Integration — Copilot with Mock LLM', function () { /** * Idle-session release exercised against the real Copilot SDK and a mock LLM. - * Uses a dedicated server with a short - * {@link AgentHostSessionReleaseGraceMsEnvVar} grace so the release fires - * promptly after the last subscriber drops (production defaults to 30s). Kept - * in its own suite/server so the short grace can't perturb the timing of the - * other agent host e2e suites. + * The dedicated server uses a zero residency cap and short provider-veto retry + * so release is deterministic without changing production policy. */ suite('Agent Host Provider Integration — Copilot Idle Release', function () { // Short enough that a post-unsubscribe wait reliably outlasts it, long // enough that the intra-test subscribe calls in createProviderSession don't race it. - const RELEASE_GRACE_MS = 500; + const RELEASE_RETRY_MS = 500; let server: IServerHandle; let client: TestProtocolClient; @@ -139,7 +136,10 @@ suite('Agent Host Provider Integration — Copilot Idle Release', function () { mockLlm: true, homeDir: suiteHome, userDataDir: join(suiteHome, 'user-data'), - env: { [AgentHostSessionReleaseGraceMsEnvVar]: String(RELEASE_GRACE_MS) }, + env: { + [AgentHostSessionResidencyLimitEnvVar]: '0', + [AgentHostSessionReleaseRetryMsEnvVar]: String(RELEASE_RETRY_MS), + }, mockScenarios: [{ id: DETACHED_SHELL_SCENARIO_ID, definition: { @@ -238,7 +238,7 @@ suite('Agent Host Provider Integration — Copilot Idle Release', function () { for (const channel of [buildDefaultChatUri(sessionUri), sessionUri]) { client.notify('unsubscribe', { channel }); } - await timeout(RELEASE_GRACE_MS + 1000); + await timeout(RELEASE_RETRY_MS + 1000); for (let attempt = 0; attempt < 150 && !existsSync(detachedCompletionMarker); attempt++) { await timeout(100); @@ -273,17 +273,12 @@ suite('Agent Host Provider Integration — Copilot Idle Release', function () { const before = await fetchSessionWithChat(client, sessionUri); assert.match(assistantMarkdown(before.turns, 'turn-release-1'), new RegExp(`\\b${firstProbe}\\b`, 'i'), 'first turn should have completed before release'); - // Drop every subscriber. The parent-session unsubscribe is sent last so it - // arms idle-session eviction on the server; after the short release grace - // elapses the cached protocol state is dropped AND the provider releases - // the live SDK session (session.disconnect), while the on-disk session log - // is preserved. + // Drop every subscriber. The zero-capacity server drops cached protocol + // state and releases the live SDK session while preserving its event log. for (const channel of [buildDefaultChatUri(sessionUri), sessionUri]) { client.notify('unsubscribe', { channel }); } - // Wait comfortably past the release grace so the release actually fires - // (and its sequenced SDK disconnect completes) before we re-subscribe. - await timeout(RELEASE_GRACE_MS + 2000); + await timeout(RELEASE_RETRY_MS + 2000); // Re-subscribe: the server restores the session from disk and the provider // resumes the SDK session on demand. The restored transcript must match