From 7c3728cd02406f79a2cdf530412a90c719104c86 Mon Sep 17 00:00:00 2001 From: Charles Vien Date: Wed, 1 Jul 2026 16:40:56 -0700 Subject: [PATCH 01/20] cap in-memory session events at 5000 --- .../src/sessions/cloudRunIdleTracker.test.ts | 36 +++++++ .../core/src/sessions/cloudRunIdleTracker.ts | 30 ++++-- .../core/src/sessions/sessionStore.test.ts | 94 +++++++++++++++++++ packages/core/src/sessions/sessionStore.ts | 12 +++ packages/shared/src/sessions.ts | 1 + 5 files changed, 166 insertions(+), 7 deletions(-) create mode 100644 packages/core/src/sessions/sessionStore.test.ts diff --git a/packages/core/src/sessions/cloudRunIdleTracker.test.ts b/packages/core/src/sessions/cloudRunIdleTracker.test.ts index 49c7ae2c13..abeb4ee135 100644 --- a/packages/core/src/sessions/cloudRunIdleTracker.test.ts +++ b/packages/core/src/sessions/cloudRunIdleTracker.test.ts @@ -130,3 +130,39 @@ describe("CloudRunIdleTracker mark/capture/restore", () => { expect(tracker.evaluateIdle(session("r1", [])).idle).toBe(false); }); }); + +describe("CloudRunIdleTracker with trimmed event history", () => { + it("keeps incremental cursors stable when the event head is trimmed", () => { + const tracker = new CloudRunIdleTracker(); + const s = session("r1", [runStarted("r1"), promptRequest()]); + expect(tracker.evaluateIdle(s).idle).toBe(false); + + const trimmed = session("r1", [turnComplete()]); + trimmed.trimmedEventCount = 2; + expect(tracker.evaluateIdle(trimmed).idle).toBe(true); + }); + + it("does not skip events appended in the same batch as a trim", () => { + const tracker = new CloudRunIdleTracker(); + const s = session("r1", [runStarted("r1")]); + expect(tracker.evaluateIdle(s).idle).toBe(false); + + const trimmed = session("r1", [promptRequest(), turnComplete()]); + trimmed.trimmedEventCount = 1; + expect(tracker.evaluateIdle(trimmed).idle).toBe(true); + }); + + it("restoreAfterFailedSend translates snapshot offsets across trims", () => { + const tracker = new CloudRunIdleTracker(); + const before = session("r1", [runStarted("r1")], "r1"); + tracker.markIdle(before); + const snapshot = tracker.capture(before); + + const after = session("r1", [], "r1"); + after.trimmedEventCount = 1; + tracker.markBusy(after); + expect(tracker.restoreAfterFailedSend(snapshot, after)).toEqual({ + agentIdleForRunId: "r1", + }); + }); +}); diff --git a/packages/core/src/sessions/cloudRunIdleTracker.ts b/packages/core/src/sessions/cloudRunIdleTracker.ts index 24ef58e367..8afac6c159 100644 --- a/packages/core/src/sessions/cloudRunIdleTracker.ts +++ b/packages/core/src/sessions/cloudRunIdleTracker.ts @@ -24,6 +24,14 @@ export interface CloudRunIdleScanResult { shouldCacheToStore: boolean; } +function streamLength(session: AgentSession): number { + return (session.trimmedEventCount ?? 0) + session.events.length; +} + +function toArrayIndex(streamOffset: number, session: AgentSession): number { + return Math.max(0, streamOffset - (session.trimmedEventCount ?? 0)); +} + /** * Tracks idleness for cloud runs incrementally so repeated `in_progress` * updates don't re-scan the full event list each time. @@ -46,7 +54,7 @@ export class CloudRunIdleTracker { */ markBusy(session: AgentSession): void { this.scanStates.set(session.taskRunId, { - nextEventIndex: session.events.length, + nextEventIndex: streamLength(session), seenCurrentRunStart: true, idle: false, }); @@ -54,7 +62,7 @@ export class CloudRunIdleTracker { markIdle(session: AgentSession): void { this.scanStates.set(session.taskRunId, { - nextEventIndex: session.events.length, + nextEventIndex: streamLength(session), seenCurrentRunStart: true, idle: true, }); @@ -64,7 +72,7 @@ export class CloudRunIdleTracker { const scanState = this.scanStates.get(session.taskRunId); return { taskRunId: session.taskRunId, - eventCount: session.events.length, + eventCount: streamLength(session), agentIdleForRunId: session.agentIdleForRunId, scanState: scanState ? { ...scanState } : undefined, }; @@ -74,7 +82,11 @@ export class CloudRunIdleTracker { snapshot: CloudRunIdleEvidenceSnapshot, session: AgentSession, ): CloudRunIdleRestoreResult | undefined { - for (let i = snapshot.eventCount; i < session.events.length; i += 1) { + for ( + let i = toArrayIndex(snapshot.eventCount, session); + i < session.events.length; + i += 1 + ) { const acpMsg = session.events[i]; if ( acpMsg && @@ -114,7 +126,7 @@ export class CloudRunIdleTracker { } let scanState = this.scanStates.get(session.taskRunId); - if (!scanState || scanState.nextEventIndex > session.events.length) { + if (!scanState || scanState.nextEventIndex > streamLength(session)) { scanState = { nextEventIndex: 0, seenCurrentRunStart: false, @@ -122,7 +134,11 @@ export class CloudRunIdleTracker { }; } - for (let i = scanState.nextEventIndex; i < session.events.length; i += 1) { + for ( + let i = toArrayIndex(scanState.nextEventIndex, session); + i < session.events.length; + i += 1 + ) { const acpMsg = session.events[i]; if (!acpMsg) continue; const msg = acpMsg.message; @@ -152,7 +168,7 @@ export class CloudRunIdleTracker { } } - scanState.nextEventIndex = session.events.length; + scanState.nextEventIndex = streamLength(session); this.scanStates.set(session.taskRunId, scanState); return { idle: scanState.idle, shouldCacheToStore: scanState.idle }; diff --git a/packages/core/src/sessions/sessionStore.test.ts b/packages/core/src/sessions/sessionStore.test.ts new file mode 100644 index 0000000000..1fb1c70db0 --- /dev/null +++ b/packages/core/src/sessions/sessionStore.test.ts @@ -0,0 +1,94 @@ +import type { AcpMessage, AgentSession } from "@posthog/shared"; +import { beforeEach, describe, expect, it } from "vitest"; +import { + MAX_SESSION_EVENTS, + sessionStore, + sessionStoreSetters, +} from "./sessionStore"; + +function event(n: number): AcpMessage { + return { + type: "acp_message", + ts: n, + message: { + jsonrpc: "2.0", + method: "session/update", + params: { sessionUpdate: "agent_message_chunk", n }, + }, + } as AcpMessage; +} + +function makeSession(taskRunId: string, taskId: string): AgentSession { + return { + taskRunId, + taskId, + events: [], + optimisticItems: [], + messageQueue: [], + pendingPermissions: new Map(), + } as unknown as AgentSession; +} + +function events(count: number, offset = 0): AcpMessage[] { + return Array.from({ length: count }, (_, i) => event(offset + i)); +} + +describe("sessionStore event cap", () => { + beforeEach(() => { + sessionStoreSetters.clearAll(); + sessionStoreSetters.setSession(makeSession("run-1", "task-1")); + }); + + it("keeps events unbounded-free: appends under the cap verbatim", () => { + sessionStoreSetters.appendEvents("run-1", events(10)); + + const session = sessionStore.getState().sessions["run-1"]; + expect(session.events).toHaveLength(10); + expect(session.trimmedEventCount ?? 0).toBe(0); + }); + + it("drops the oldest events beyond MAX_SESSION_EVENTS and records the trim", () => { + sessionStoreSetters.appendEvents("run-1", events(MAX_SESSION_EVENTS + 100)); + + const session = sessionStore.getState().sessions["run-1"]; + expect(session.events).toHaveLength(MAX_SESSION_EVENTS); + expect(session.trimmedEventCount).toBe(100); + expect( + (session.events[0].message as { params: { n: number } }).params.n, + ).toBe(100); + }); + + it("accumulates trimmedEventCount across appends", () => { + sessionStoreSetters.appendEvents("run-1", events(MAX_SESSION_EVENTS)); + sessionStoreSetters.appendEvents("run-1", events(50, MAX_SESSION_EVENTS)); + sessionStoreSetters.appendEvents( + "run-1", + events(25, MAX_SESSION_EVENTS + 50), + ); + + const session = sessionStore.getState().sessions["run-1"]; + expect(session.events).toHaveLength(MAX_SESSION_EVENTS); + expect(session.trimmedEventCount).toBe(75); + expect( + (session.events.at(-1)?.message as { params: { n: number } }).params.n, + ).toBe(MAX_SESSION_EVENTS + 74); + }); + + it("trims in replaceOptimisticWithEvent too", () => { + sessionStoreSetters.appendEvents("run-1", events(MAX_SESSION_EVENTS)); + sessionStoreSetters.appendOptimisticItem("run-1", { + type: "user_message", + content: "hi", + } as never); + + sessionStoreSetters.replaceOptimisticWithEvent( + "run-1", + event(MAX_SESSION_EVENTS), + ); + + const session = sessionStore.getState().sessions["run-1"]; + expect(session.events).toHaveLength(MAX_SESSION_EVENTS); + expect(session.trimmedEventCount).toBe(1); + expect(session.optimisticItems).toHaveLength(0); + }); +}); diff --git a/packages/core/src/sessions/sessionStore.ts b/packages/core/src/sessions/sessionStore.ts index 0e86a78c6a..82c24b622a 100644 --- a/packages/core/src/sessions/sessionStore.ts +++ b/packages/core/src/sessions/sessionStore.ts @@ -32,6 +32,16 @@ export const sessionStore = createStore()( })), ); +export const MAX_SESSION_EVENTS = 5000; + +function trimSessionEvents(session: AgentSession) { + const excess = session.events.length - MAX_SESSION_EVENTS; + if (excess > 0) { + session.events.splice(0, excess); + session.trimmedEventCount = (session.trimmedEventCount ?? 0) + excess; + } +} + export const sessionStoreSetters = { setSession: (session: AgentSession) => { sessionStore.setState((state) => { @@ -76,6 +86,7 @@ export const sessionStoreSetters = { // immer autofreeze, so this is the only freeze. for (const event of events) Object.freeze(event); session.events.push(...events); + trimSessionEvents(session); if (newLineCount !== undefined) { session.processedLineCount = newLineCount; } @@ -301,6 +312,7 @@ export const sessionStoreSetters = { const session = state.sessions[taskRunId]; if (session) { session.events.push(Object.freeze(event)); + trimSessionEvents(session); session.optimisticItems = []; } }); diff --git a/packages/shared/src/sessions.ts b/packages/shared/src/sessions.ts index 0724dddfac..41728de954 100644 --- a/packages/shared/src/sessions.ts +++ b/packages/shared/src/sessions.ts @@ -62,6 +62,7 @@ export interface AgentSession { currentPromptId?: number | null; logUrl?: string; processedLineCount?: number; + trimmedEventCount?: number; framework?: "claude"; adapter?: Adapter; configOptions?: SessionConfigOption[]; From e746f25623ddc1fea29412575a258c9e65111341 Mon Sep 17 00:00:00 2001 From: Charles Vien Date: Wed, 1 Jul 2026 16:41:11 -0700 Subject: [PATCH 02/20] evict idle sessions beyond an LRU budget --- .../core/src/sessions/sessionEviction.test.ts | 103 ++++++++++++++++++ packages/core/src/sessions/sessionEviction.ts | 33 ++++++ packages/core/src/sessions/sessionService.ts | 29 +++++ 3 files changed, 165 insertions(+) create mode 100644 packages/core/src/sessions/sessionEviction.test.ts create mode 100644 packages/core/src/sessions/sessionEviction.ts diff --git a/packages/core/src/sessions/sessionEviction.test.ts b/packages/core/src/sessions/sessionEviction.test.ts new file mode 100644 index 0000000000..d22987718e --- /dev/null +++ b/packages/core/src/sessions/sessionEviction.test.ts @@ -0,0 +1,103 @@ +import type { AgentSession } from "@posthog/shared"; +import { describe, expect, it } from "vitest"; +import { isSessionIdle, selectSessionsToEvict } from "./sessionEviction"; + +function makeSession(overrides: Partial): AgentSession { + return { + taskRunId: `run-${overrides.taskId}`, + status: "connected", + isPromptPending: false, + pendingPermissions: new Map(), + messageQueue: [], + startedAt: 0, + ...overrides, + } as AgentSession; +} + +describe("isSessionIdle", () => { + it.each([ + ["connected idle local session", {}, true], + ["connecting session", { status: "connecting" as const }, false], + ["pending prompt", { isPromptPending: true }, false], + [ + "pending permission", + { pendingPermissions: new Map([["p1", {} as never]]) }, + false, + ], + [ + "queued messages", + { messageQueue: [{ id: "m1", content: "x", queuedAt: 0 }] }, + false, + ], + [ + "running cloud session", + { isCloud: true, cloudStatus: "in_progress" as const }, + false, + ], + [ + "queued cloud session", + { isCloud: true, cloudStatus: "queued" as const }, + false, + ], + [ + "completed cloud session", + { isCloud: true, cloudStatus: "completed" as const }, + true, + ], + ["cloud session without status", { isCloud: true }, false], + ])("%s -> %s", (_name, overrides, expected) => { + expect(isSessionIdle(makeSession({ taskId: "t", ...overrides }))).toBe( + expected, + ); + }); +}); + +describe("selectSessionsToEvict", () => { + const lastUsedAt = (session: AgentSession) => session.startedAt; + + it("returns nothing under the budget", () => { + const sessions = [ + makeSession({ taskId: "a" }), + makeSession({ taskId: "b" }), + ]; + expect( + selectSessionsToEvict({ + sessions, + activeTaskId: "a", + lastUsedAt, + maxSessions: 3, + }), + ).toEqual([]); + }); + + it("evicts the least recently used idle sessions over the budget", () => { + const sessions = [ + makeSession({ taskId: "a", startedAt: 30 }), + makeSession({ taskId: "b", startedAt: 10 }), + makeSession({ taskId: "c", startedAt: 20 }), + makeSession({ taskId: "d", startedAt: 40 }), + ]; + const evicted = selectSessionsToEvict({ + sessions, + activeTaskId: "d", + lastUsedAt, + maxSessions: 3, + }); + expect(evicted.map((s) => s.taskId)).toEqual(["b", "c"]); + }); + + it("never evicts the active task or busy sessions", () => { + const sessions = [ + makeSession({ taskId: "active", startedAt: 1 }), + makeSession({ taskId: "busy", startedAt: 2, isPromptPending: true }), + makeSession({ taskId: "idle", startedAt: 3 }), + ]; + const evicted = selectSessionsToEvict({ + sessions, + activeTaskId: "active", + lastUsedAt, + maxSessions: 2, + }); + expect(evicted.map((s) => s.taskId)).toEqual(["idle"]); + }); +}); diff --git a/packages/core/src/sessions/sessionEviction.ts b/packages/core/src/sessions/sessionEviction.ts new file mode 100644 index 0000000000..9299ce772c --- /dev/null +++ b/packages/core/src/sessions/sessionEviction.ts @@ -0,0 +1,33 @@ +import type { AgentSession } from "@posthog/shared"; +import { isTerminalStatus } from "@posthog/shared/domain-types"; + +export const MAX_CONNECTED_SESSIONS = 8; + +export function isSessionIdle(session: AgentSession): boolean { + if (session.status === "connecting") return false; + if (session.isPromptPending) return false; + if (session.pendingPermissions.size > 0) return false; + if (session.messageQueue.length > 0) return false; + if (session.isCloud) return isTerminalStatus(session.cloudStatus); + return true; +} + +export function selectSessionsToEvict(params: { + sessions: AgentSession[]; + activeTaskId: string; + lastUsedAt: (session: AgentSession) => number; + maxSessions?: number; +}): AgentSession[] { + const { sessions, activeTaskId, lastUsedAt } = params; + const maxSessions = params.maxSessions ?? MAX_CONNECTED_SESSIONS; + + const excess = sessions.length - (maxSessions - 1); + if (excess <= 0) return []; + + return sessions + .filter( + (session) => session.taskId !== activeTaskId && isSessionIdle(session), + ) + .sort((a, b) => lastUsedAt(a) - lastUsedAt(b)) + .slice(0, excess); +} diff --git a/packages/core/src/sessions/sessionService.ts b/packages/core/src/sessions/sessionService.ts index 5b5af8740e..2423cabbf4 100644 --- a/packages/core/src/sessions/sessionService.ts +++ b/packages/core/src/sessions/sessionService.ts @@ -77,6 +77,7 @@ import { promptReferencesAbsoluteFolder, shellExecutesToContextBlocks, } from "./sessionEvents"; +import { selectSessionsToEvict } from "./sessionEviction"; import { createBaseSession } from "./sessionFactory"; import { type ParsedSessionLogs, parseSessionLogContent } from "./sessionLogs"; @@ -513,6 +514,7 @@ export class SessionService { >(); private localRepoPaths = new Map(); private localRecoveryAttempts = new Map>(); + private sessionLastUsedAt = new Map(); /** Re-entrance guard for cloud queue dispatch (per taskId). */ private dispatchingCloudQueues = new Set(); /** Coalesces deferred cloud queue flush timers (per taskId). */ @@ -610,6 +612,8 @@ export class SessionService { const { task } = params; const taskId = task.id; this.localRepoPaths.set(taskId, params.repoPath); + this.sessionLastUsedAt.set(taskId, Date.now()); + void this.evictIdleSessions(taskId); // Return existing connection promise if already connecting const existingPromise = this.connectingTasks.get(taskId); @@ -1358,6 +1362,31 @@ export class SessionService { await this.teardownSession(session.taskRunId); } + private async evictIdleSessions(activeTaskId: string): Promise { + const toEvict = selectSessionsToEvict({ + sessions: Object.values(this.d.store.getSessions()), + activeTaskId, + lastUsedAt: (session) => + this.sessionLastUsedAt.get(session.taskId) ?? session.startedAt, + }); + + for (const session of toEvict) { + this.d.log.info("Evicting idle session to bound memory", { + taskId: session.taskId, + taskRunId: session.taskRunId, + }); + this.sessionLastUsedAt.delete(session.taskId); + try { + await this.teardownSession(session.taskRunId); + } catch (error) { + this.d.log.error("Failed to evict idle session", { + taskId: session.taskId, + error, + }); + } + } + } + // --- Subscription Management --- /** Streamed events awaiting their frame flush, keyed by taskRunId. Order From 7283a1f93f2bd5079f7d0be26b0913222ba4c049 Mon Sep 17 00:00:00 2001 From: Charles Vien Date: Wed, 1 Jul 2026 16:41:25 -0700 Subject: [PATCH 03/20] destroy task terminals on suspend and archive --- .../ui/src/features/archive/useArchiveTask.ts | 7 +++++-- .../src/features/suspension/useSuspendTask.ts | 2 ++ .../src/features/terminal/TerminalManager.ts | 20 +++++++++++-------- 3 files changed, 19 insertions(+), 10 deletions(-) diff --git a/packages/ui/src/features/archive/useArchiveTask.ts b/packages/ui/src/features/archive/useArchiveTask.ts index 16901cc016..b214ef8727 100644 --- a/packages/ui/src/features/archive/useArchiveTask.ts +++ b/packages/ui/src/features/archive/useArchiveTask.ts @@ -19,6 +19,7 @@ import { useHostTRPC } from "@posthog/host-router/react"; import { useCommandCenterStore } from "@posthog/ui/features/command-center/commandCenterStore"; import { useFocusStore } from "@posthog/ui/features/focus/focusStore"; import { pinnedTasksApi } from "@posthog/ui/features/sidebar/taskMetaApi"; +import { terminalManager } from "@posthog/ui/features/terminal/TerminalManager"; import { type TerminalState, useTerminalStore, @@ -102,8 +103,10 @@ function makeOrchestrationDeps( ([key]) => key === taskId || key.startsWith(`${taskId}-`), ), ), - clearTerminalStates: (taskId) => - useTerminalStore.getState().clearTerminalStatesForTask(taskId), + clearTerminalStates: (taskId) => { + terminalManager.destroyForTask(taskId); + useTerminalStore.getState().clearTerminalStatesForTask(taskId); + }, restoreTerminalStates: (states) => { useTerminalStore.setState((s) => ({ terminalStates: { diff --git a/packages/ui/src/features/suspension/useSuspendTask.ts b/packages/ui/src/features/suspension/useSuspendTask.ts index d5482ecd80..c9e9b0c7b8 100644 --- a/packages/ui/src/features/suspension/useSuspendTask.ts +++ b/packages/ui/src/features/suspension/useSuspendTask.ts @@ -1,5 +1,6 @@ import { useHostTRPC, useHostTRPCClient } from "@posthog/host-router/react"; import { useFocusStore } from "@posthog/ui/features/focus/focusStore"; +import { terminalManager } from "@posthog/ui/features/terminal/TerminalManager"; import { useTerminalStore } from "@posthog/ui/features/terminal/terminalStore"; import { logger } from "@posthog/ui/shell/logger"; import { useMutation, useQueryClient } from "@tanstack/react-query"; @@ -29,6 +30,7 @@ export function useSuspendTask() { const workspaces = await hostClient.workspace.getAll.query(); const workspace = workspaces[taskId] ?? null; + terminalManager.destroyForTask(taskId); useTerminalStore.getState().clearTerminalStatesForTask(taskId); queryClient.setQueryData(suspendedTaskIdsKey, (old) => diff --git a/packages/ui/src/features/terminal/TerminalManager.ts b/packages/ui/src/features/terminal/TerminalManager.ts index d381263d3f..f91ed19556 100644 --- a/packages/ui/src/features/terminal/TerminalManager.ts +++ b/packages/ui/src/features/terminal/TerminalManager.ts @@ -565,9 +565,21 @@ class TerminalManagerImpl { instance.term.dispose(); + instance.terminalElement?.remove(); + instance.terminalElement = null; + this.instances.delete(sessionId); } + destroyForTask(taskId: string): void { + for (const [sessionId, instance] of this.instances) { + const key = instance.persistenceKey; + if (key === taskId || key.startsWith(`${taskId}-`)) { + this.destroy(sessionId); + } + } + } + focus(sessionId: string): void { const instance = this.instances.get(sessionId); if (instance) { @@ -675,14 +687,6 @@ class TerminalManagerImpl { } } - destroyByPrefix(prefix: string): void { - for (const sessionId of this.instances.keys()) { - if (sessionId.startsWith(prefix)) { - this.destroy(sessionId); - } - } - } - getSessionsByPrefix(prefix: string): string[] { const result: string[] = []; for (const sessionId of this.instances.keys()) { From b8ebd0b0849d70f0d413c73513c30c98deedae5d Mon Sep 17 00:00:00 2001 From: Charles Vien Date: Wed, 1 Jul 2026 16:41:27 -0700 Subject: [PATCH 04/20] release session resources when deleting a task --- .../features/tasks/useTaskCrudMutations.ts | 28 ++++++++++++++++++- 1 file changed, 27 insertions(+), 1 deletion(-) diff --git a/packages/ui/src/features/tasks/useTaskCrudMutations.ts b/packages/ui/src/features/tasks/useTaskCrudMutations.ts index 70111964e1..e3f13044b1 100644 --- a/packages/ui/src/features/tasks/useTaskCrudMutations.ts +++ b/packages/ui/src/features/tasks/useTaskCrudMutations.ts @@ -1,3 +1,7 @@ +import { + SESSION_SERVICE, + type SessionService, +} from "@posthog/core/sessions/sessionService"; import { insertTaskDedup, removeTaskFromList, @@ -6,13 +10,31 @@ import { TASK_DELETION_SERVICE, type TaskDeletionService, } from "@posthog/core/tasks/taskDeletionService"; +import { resolveService } from "@posthog/di/container"; import { useService } from "@posthog/di/react"; import type { Task } from "@posthog/shared/domain-types"; +import { terminalManager } from "@posthog/ui/features/terminal/TerminalManager"; +import { useTerminalStore } from "@posthog/ui/features/terminal/terminalStore"; import { useAuthenticatedMutation } from "@posthog/ui/hooks/useAuthenticatedMutation"; +import { logger } from "@posthog/ui/shell/logger"; import { useQueryClient } from "@tanstack/react-query"; import { useCallback } from "react"; import { taskKeys } from "./taskKeys"; +const log = logger.scope("tasks"); + +async function releaseDeletedTaskResources(taskId: string): Promise { + try { + await resolveService(SESSION_SERVICE).disconnectFromTask( + taskId, + ); + } catch (error) { + log.error("Failed to disconnect session for deleted task", error); + } + terminalManager.destroyForTask(taskId); + useTerminalStore.getState().clearTerminalStatesForTask(taskId); +} + export function useCreateTask() { const queryClient = useQueryClient(); @@ -77,7 +99,11 @@ export function useDeleteTask() { ); const mutation = useAuthenticatedMutation( - (client, taskId: string) => deletionService.deleteTask(client, taskId), + async (client, taskId: string) => { + const result = await deletionService.deleteTask(client, taskId); + await releaseDeletedTaskResources(taskId); + return result; + }, { onMutate: async (taskId) => { await queryClient.cancelQueries({ queryKey: taskKeys.lists() }); From f53b92b409c4311818e69a263afc09d4f96b0640 Mon Sep 17 00:00:00 2001 From: Charles Vien Date: Wed, 1 Jul 2026 16:41:28 -0700 Subject: [PATCH 05/20] slow the sidebar inbox badge poll to 60s --- .../ui/src/features/inbox/hooks/useInboxAllReports.ts | 7 +++++-- .../features/sidebar/components/SidebarNavSection.tsx | 9 ++++++++- 2 files changed, 13 insertions(+), 3 deletions(-) diff --git a/packages/ui/src/features/inbox/hooks/useInboxAllReports.ts b/packages/ui/src/features/inbox/hooks/useInboxAllReports.ts index 7e8115e84a..99f19b7a33 100644 --- a/packages/ui/src/features/inbox/hooks/useInboxAllReports.ts +++ b/packages/ui/src/features/inbox/hooks/useInboxAllReports.ts @@ -42,9 +42,12 @@ export function useInboxAllReports(options?: { ignoreScope?: boolean; ignoreFilters?: boolean; pullRequestsOnly?: boolean; + refetchIntervalMs?: number; }) { const ignoreScope = options?.ignoreScope ?? false; const ignoreFilters = options?.ignoreFilters ?? false; + const refetchIntervalMs = + options?.refetchIntervalMs ?? INBOX_REFETCH_INTERVAL_MS; // The Pull requests tab fetches a server-filtered list (reports that have a // shipped PR) so its list body comes from the same source as its count — a PR // sitting past the broad list's first page no longer renders an empty tab @@ -101,7 +104,7 @@ export function useInboxAllReports(options?: { // throwaway project-wide fetch first. Other scopes don't depend on the // user and run immediately. enabled: !isForYou || reviewerUuid != null, - refetchInterval: INBOX_REFETCH_INTERVAL_MS, + refetchInterval: refetchIntervalMs, refetchIntervalInBackground: false, }, ); @@ -130,7 +133,7 @@ export function useInboxAllReports(options?: { }, { enabled: !isForYou || reviewerUuid != null, - refetchInterval: INBOX_REFETCH_INTERVAL_MS, + refetchInterval: refetchIntervalMs, refetchIntervalInBackground: false, }, ); diff --git a/packages/ui/src/features/sidebar/components/SidebarNavSection.tsx b/packages/ui/src/features/sidebar/components/SidebarNavSection.tsx index b13a01201a..211d722d96 100644 --- a/packages/ui/src/features/sidebar/components/SidebarNavSection.tsx +++ b/packages/ui/src/features/sidebar/components/SidebarNavSection.tsx @@ -30,6 +30,8 @@ import { NewTaskItem } from "./items/NewTaskItem"; import { SearchItem } from "./items/SearchItem"; import { SkillsItem } from "./items/SkillsItem"; +const SIDEBAR_INBOX_REFETCH_INTERVAL_MS = 60_000; + interface SidebarNavSectionProps { // The Command Center badge counts how many command-center cells point at a // live task. Deriving it needs the task list, which the Code pane's @@ -85,7 +87,12 @@ export function SidebarNavSection({ // Pull requests tab shows, so the badge and the tab always agree. // `ignoreFilters` keeps the badge stable against the inbox's filter chrome; // scope still follows the user's For-you / project choice. - const { counts: inboxCounts } = useInboxAllReports({ ignoreFilters: true }); + // The sidebar mounts on every route, so its badge polls slowly; opening the + // inbox adds its own 3s observers and React Query uses the shortest interval. + const { counts: inboxCounts } = useInboxAllReports({ + ignoreFilters: true, + refetchIntervalMs: SIDEBAR_INBOX_REFETCH_INTERVAL_MS, + }); const inboxPullRequestCount = inboxCounts.pulls; // Only subscribe to the task list when a parent hasn't already supplied the From 60595ead6f1558a9e6973cd5041bd1c42e37182f Mon Sep 17 00:00:00 2001 From: Charles Vien Date: Wed, 1 Jul 2026 16:41:30 -0700 Subject: [PATCH 06/20] relax connectivity polling while online --- .../src/services/connectivity/service.test.ts | 5 ++-- .../src/services/connectivity/service.ts | 23 ++++++++++--------- 2 files changed, 15 insertions(+), 13 deletions(-) diff --git a/packages/workspace-server/src/services/connectivity/service.test.ts b/packages/workspace-server/src/services/connectivity/service.test.ts index ee135193ee..44be28df52 100644 --- a/packages/workspace-server/src/services/connectivity/service.test.ts +++ b/packages/workspace-server/src/services/connectivity/service.test.ts @@ -104,7 +104,8 @@ describe("ConnectivityService", () => { service.on(ConnectivityEvent.StatusChange, handler); mockFetch.mockRejectedValue(new Error("offline")); - await vi.advanceTimersByTimeAsync(6000); // two failed polls + await vi.advanceTimersByTimeAsync(30_000); // 1st failure at the healthy cadence + await vi.advanceTimersByTimeAsync(3000); // fast recheck confirms offline expect(handler).toHaveBeenCalledWith({ isOnline: false }); expect(handler).toHaveBeenCalledTimes(1); @@ -219,7 +220,7 @@ describe("ConnectivityService", () => { const callsAfterInit = mockFetch.mock.calls.length; - await vi.advanceTimersByTimeAsync(3000); + await vi.advanceTimersByTimeAsync(30_000); expect(mockFetch.mock.calls.length).toBeGreaterThan(callsAfterInit); }); }); diff --git a/packages/workspace-server/src/services/connectivity/service.ts b/packages/workspace-server/src/services/connectivity/service.ts index fed5859bbe..114c25fbe2 100644 --- a/packages/workspace-server/src/services/connectivity/service.ts +++ b/packages/workspace-server/src/services/connectivity/service.ts @@ -14,7 +14,7 @@ const CHECK_TIMEOUT_MS = 5_000; const OFFLINE_CONFIRM_THRESHOLD = 2; const MIN_POLL_INTERVAL_MS = 3_000; const MAX_POLL_INTERVAL_MS = 10_000; -const ONLINE_POLL_INTERVAL_MS = 3_000; +const ONLINE_POLL_INTERVAL_MS = 30_000; const OFFLINE_BACKOFF_MULTIPLIER = 1.5; @injectable() @@ -27,8 +27,7 @@ export class ConnectivityService extends TypedEventEmitter { constructor() { super(); this.setMaxListeners(0); - void this.checkConnectivity(); - this.startPolling(); + void this.checkConnectivity().finally(() => this.startPolling()); } getStatus(): ConnectivityStatusOutput { @@ -76,14 +75,13 @@ export class ConnectivityService extends TypedEventEmitter { } private async verifyWithHttp(): Promise { - try { - // Resolves as soon as the first host responds reachably; rejects only - // when every host fails. - await Promise.any(CHECK_URLS.map((url) => this.probe(url))); - return true; - } catch { - return false; + for (const url of CHECK_URLS) { + try { + await this.probe(url); + return true; + } catch {} } + return false; } private async probe(url: string): Promise { @@ -103,8 +101,11 @@ export class ConnectivityService extends TypedEventEmitter { } private schedulePoll(): void { + // Poll rarely while healthy, quickly while confirming a suspected outage. const interval = this.isOnline - ? ONLINE_POLL_INTERVAL_MS + ? this.consecutiveFailures > 0 + ? MIN_POLL_INTERVAL_MS + : ONLINE_POLL_INTERVAL_MS : Math.min( MIN_POLL_INTERVAL_MS * OFFLINE_BACKOFF_MULTIPLIER ** this.offlinePollAttempt, From 4ce63030bdd29a333e8b08193a6b1d431008e9b8 Mon Sep 17 00:00:00 2001 From: Charles Vien Date: Wed, 1 Jul 2026 17:41:15 -0700 Subject: [PATCH 07/20] trim hydrated session histories to the event cap --- .../core/src/sessions/sessionStore.test.ts | 29 +++++++++++++++++++ packages/core/src/sessions/sessionStore.ts | 4 +++ 2 files changed, 33 insertions(+) diff --git a/packages/core/src/sessions/sessionStore.test.ts b/packages/core/src/sessions/sessionStore.test.ts index 1fb1c70db0..5d565b93eb 100644 --- a/packages/core/src/sessions/sessionStore.test.ts +++ b/packages/core/src/sessions/sessionStore.test.ts @@ -91,4 +91,33 @@ describe("sessionStore event cap", () => { expect(session.trimmedEventCount).toBe(1); expect(session.optimisticItems).toHaveLength(0); }); + + it("trims oversized histories installed via setSession", () => { + const hydrated = makeSession("run-2", "task-2"); + hydrated.events = events(MAX_SESSION_EVENTS + 200); + sessionStoreSetters.setSession(hydrated); + + const session = sessionStore.getState().sessions["run-2"]; + expect(session.events).toHaveLength(MAX_SESSION_EVENTS); + expect(session.trimmedEventCount).toBe(200); + }); + + it("trims oversized histories installed via updateSession", () => { + sessionStoreSetters.updateSession("run-1", { + events: events(MAX_SESSION_EVENTS + 30), + }); + + const session = sessionStore.getState().sessions["run-1"]; + expect(session.events).toHaveLength(MAX_SESSION_EVENTS); + expect(session.trimmedEventCount).toBe(30); + }); + + it("leaves untouched events alone in updateSession without events", () => { + sessionStoreSetters.appendEvents("run-1", events(10)); + sessionStoreSetters.updateSession("run-1", { taskTitle: "renamed" }); + + const session = sessionStore.getState().sessions["run-1"]; + expect(session.events).toHaveLength(10); + expect(session.trimmedEventCount ?? 0).toBe(0); + }); }); diff --git a/packages/core/src/sessions/sessionStore.ts b/packages/core/src/sessions/sessionStore.ts index 82c24b622a..8054740500 100644 --- a/packages/core/src/sessions/sessionStore.ts +++ b/packages/core/src/sessions/sessionStore.ts @@ -53,6 +53,7 @@ export const sessionStoreSetters = { state.sessions[session.taskRunId] = session; state.taskIdIndex[session.taskId] = session.taskRunId; + trimSessionEvents(state.sessions[session.taskRunId]); }); }, @@ -70,6 +71,9 @@ export const sessionStoreSetters = { sessionStore.setState((state) => { if (state.sessions[taskRunId]) { Object.assign(state.sessions[taskRunId], updates); + if (updates.events) { + trimSessionEvents(state.sessions[taskRunId]); + } } }); }, From 6f20787c15cee05f9c77af9cb8af184fb396e123 Mon Sep 17 00:00:00 2001 From: Charles Vien Date: Wed, 1 Jul 2026 17:41:31 -0700 Subject: [PATCH 08/20] protect mounted tasks from session eviction --- .../core/src/sessions/sessionEviction.test.ts | 96 ++++++------ packages/core/src/sessions/sessionEviction.ts | 11 +- packages/core/src/sessions/sessionService.ts | 21 +++ .../sessions/sessionServiceEviction.test.ts | 141 ++++++++++++++++++ .../sessions/hooks/useSessionConnection.ts | 4 + 5 files changed, 228 insertions(+), 45 deletions(-) create mode 100644 packages/core/src/sessions/sessionServiceEviction.test.ts diff --git a/packages/core/src/sessions/sessionEviction.test.ts b/packages/core/src/sessions/sessionEviction.test.ts index d22987718e..26fb8ee228 100644 --- a/packages/core/src/sessions/sessionEviction.test.ts +++ b/packages/core/src/sessions/sessionEviction.test.ts @@ -45,6 +45,8 @@ describe("isSessionIdle", () => { true, ], ["cloud session without status", { isCloud: true }, false], + ["disconnected local session", { status: "disconnected" as const }, true], + ["errored local session", { status: "error" as const }, true], ])("%s -> %s", (_name, overrides, expected) => { expect(isSessionIdle(makeSession({ taskId: "t", ...overrides }))).toBe( expected, @@ -55,49 +57,59 @@ describe("isSessionIdle", () => { describe("selectSessionsToEvict", () => { const lastUsedAt = (session: AgentSession) => session.startedAt; - it("returns nothing under the budget", () => { - const sessions = [ - makeSession({ taskId: "a" }), - makeSession({ taskId: "b" }), - ]; - expect( - selectSessionsToEvict({ - sessions, + it.each([ + [ + "returns nothing under the budget", + { + sessions: [makeSession({ taskId: "a" }), makeSession({ taskId: "b" })], activeTaskId: "a", - lastUsedAt, maxSessions: 3, - }), - ).toEqual([]); - }); - - it("evicts the least recently used idle sessions over the budget", () => { - const sessions = [ - makeSession({ taskId: "a", startedAt: 30 }), - makeSession({ taskId: "b", startedAt: 10 }), - makeSession({ taskId: "c", startedAt: 20 }), - makeSession({ taskId: "d", startedAt: 40 }), - ]; - const evicted = selectSessionsToEvict({ - sessions, - activeTaskId: "d", - lastUsedAt, - maxSessions: 3, - }); - expect(evicted.map((s) => s.taskId)).toEqual(["b", "c"]); - }); - - it("never evicts the active task or busy sessions", () => { - const sessions = [ - makeSession({ taskId: "active", startedAt: 1 }), - makeSession({ taskId: "busy", startedAt: 2, isPromptPending: true }), - makeSession({ taskId: "idle", startedAt: 3 }), - ]; - const evicted = selectSessionsToEvict({ - sessions, - activeTaskId: "active", - lastUsedAt, - maxSessions: 2, - }); - expect(evicted.map((s) => s.taskId)).toEqual(["idle"]); + }, + [], + ], + [ + "evicts the least recently used idle sessions over the budget", + { + sessions: [ + makeSession({ taskId: "a", startedAt: 30 }), + makeSession({ taskId: "b", startedAt: 10 }), + makeSession({ taskId: "c", startedAt: 20 }), + makeSession({ taskId: "d", startedAt: 40 }), + ], + activeTaskId: "d", + maxSessions: 3, + }, + ["b", "c"], + ], + [ + "never evicts the active task or busy sessions", + { + sessions: [ + makeSession({ taskId: "active", startedAt: 1 }), + makeSession({ taskId: "busy", startedAt: 2, isPromptPending: true }), + makeSession({ taskId: "idle", startedAt: 3 }), + ], + activeTaskId: "active", + maxSessions: 2, + }, + ["idle"], + ], + [ + "never evicts mounted tasks", + { + sessions: [ + makeSession({ taskId: "a", startedAt: 1 }), + makeSession({ taskId: "b", startedAt: 2 }), + makeSession({ taskId: "c", startedAt: 3 }), + ], + activeTaskId: "c", + protectedTaskIds: new Set(["a"]), + maxSessions: 2, + }, + ["b"], + ], + ])("%s", (_name, params, expected) => { + const evicted = selectSessionsToEvict({ ...params, lastUsedAt }); + expect(evicted.map((s) => s.taskId)).toEqual(expected); }); }); diff --git a/packages/core/src/sessions/sessionEviction.ts b/packages/core/src/sessions/sessionEviction.ts index 9299ce772c..d56012d875 100644 --- a/packages/core/src/sessions/sessionEviction.ts +++ b/packages/core/src/sessions/sessionEviction.ts @@ -1,7 +1,8 @@ import type { AgentSession } from "@posthog/shared"; import { isTerminalStatus } from "@posthog/shared/domain-types"; -export const MAX_CONNECTED_SESSIONS = 8; +// Above the Command Center's 3x3 grid so fully-visible layouts never evict. +export const MAX_CONNECTED_SESSIONS = 12; export function isSessionIdle(session: AgentSession): boolean { if (session.status === "connecting") return false; @@ -15,10 +16,11 @@ export function isSessionIdle(session: AgentSession): boolean { export function selectSessionsToEvict(params: { sessions: AgentSession[]; activeTaskId: string; + protectedTaskIds?: ReadonlySet; lastUsedAt: (session: AgentSession) => number; maxSessions?: number; }): AgentSession[] { - const { sessions, activeTaskId, lastUsedAt } = params; + const { sessions, activeTaskId, protectedTaskIds, lastUsedAt } = params; const maxSessions = params.maxSessions ?? MAX_CONNECTED_SESSIONS; const excess = sessions.length - (maxSessions - 1); @@ -26,7 +28,10 @@ export function selectSessionsToEvict(params: { return sessions .filter( - (session) => session.taskId !== activeTaskId && isSessionIdle(session), + (session) => + session.taskId !== activeTaskId && + !protectedTaskIds?.has(session.taskId) && + isSessionIdle(session), ) .sort((a, b) => lastUsedAt(a) - lastUsedAt(b)) .slice(0, excess); diff --git a/packages/core/src/sessions/sessionService.ts b/packages/core/src/sessions/sessionService.ts index 2423cabbf4..30791b4be4 100644 --- a/packages/core/src/sessions/sessionService.ts +++ b/packages/core/src/sessions/sessionService.ts @@ -515,6 +515,7 @@ export class SessionService { private localRepoPaths = new Map(); private localRecoveryAttempts = new Map>(); private sessionLastUsedAt = new Map(); + private mountedTaskCounts = new Map(); /** Re-entrance guard for cloud queue dispatch (per taskId). */ private dispatchingCloudQueues = new Set(); /** Coalesces deferred cloud queue flush timers (per taskId). */ @@ -1057,6 +1058,7 @@ export class SessionService { if (session) { this.localRepoPaths.delete(session.taskId); this.localRecoveryAttempts.delete(session.taskId); + this.sessionLastUsedAt.delete(session.taskId); } this.d.adapterStore.removeAdapter(taskRunId); this.d.removePersistedConfigOptions(taskRunId); @@ -1362,10 +1364,28 @@ export class SessionService { await this.teardownSession(session.taskRunId); } + registerMountedTask(taskId: string): () => void { + this.mountedTaskCounts.set( + taskId, + (this.mountedTaskCounts.get(taskId) ?? 0) + 1, + ); + this.sessionLastUsedAt.set(taskId, Date.now()); + return () => { + const count = this.mountedTaskCounts.get(taskId) ?? 0; + if (count <= 1) { + this.mountedTaskCounts.delete(taskId); + } else { + this.mountedTaskCounts.set(taskId, count - 1); + } + this.sessionLastUsedAt.set(taskId, Date.now()); + }; + } + private async evictIdleSessions(activeTaskId: string): Promise { const toEvict = selectSessionsToEvict({ sessions: Object.values(this.d.store.getSessions()), activeTaskId, + protectedTaskIds: new Set(this.mountedTaskCounts.keys()), lastUsedAt: (session) => this.sessionLastUsedAt.get(session.taskId) ?? session.startedAt, }); @@ -1627,6 +1647,7 @@ export class SessionService { this.connectingTasks.clear(); this.localRepoPaths.clear(); this.localRecoveryAttempts.clear(); + this.sessionLastUsedAt.clear(); this.cloudPermissionRequestIds.clear(); this.liveTurnContent.clear(); this.cloudLogGapReconciler.clear(); diff --git a/packages/core/src/sessions/sessionServiceEviction.test.ts b/packages/core/src/sessions/sessionServiceEviction.test.ts new file mode 100644 index 0000000000..ce329977c2 --- /dev/null +++ b/packages/core/src/sessions/sessionServiceEviction.test.ts @@ -0,0 +1,141 @@ +import type { AgentSession } from "@posthog/shared"; +import type { Task } from "@posthog/shared/domain-types"; +import { describe, expect, it, vi } from "vitest"; +import { MAX_CONNECTED_SESSIONS } from "./sessionEviction"; +import { SessionService, type SessionServiceDeps } from "./sessionService"; + +function makeSession( + taskId: string, + startedAt: number, + overrides: Partial = {}, +): AgentSession { + return { + taskRunId: `run-${taskId}`, + taskId, + taskTitle: taskId, + channel: "", + events: [], + startedAt, + status: "connected", + isPromptPending: false, + isCompacting: false, + promptStartedAt: null, + pendingPermissions: new Map(), + pausedDurationMs: 0, + messageQueue: [], + optimisticItems: [], + ...overrides, + } as AgentSession; +} + +function createHarness(seedSessions: AgentSession[]) { + const sessions: Record = {}; + for (const session of seedSessions) { + sessions[session.taskRunId] = session; + } + const removeSession = vi.fn((taskRunId: string) => { + delete sessions[taskRunId]; + }); + const cancelMutate = vi.fn().mockResolvedValue(undefined); + + const store = { + getSessions: () => sessions, + getSessionByTaskId: (taskId: string) => + Object.values(sessions).find((s) => s.taskId === taskId), + removeSession, + }; + + const noopLog = { + info: vi.fn(), + warn: vi.fn(), + error: vi.fn(), + debug: vi.fn(), + }; + + const deps = { + store, + log: noopLog, + getPersistedConfigOptions: () => undefined, + setPersistedConfigOptions: vi.fn(), + removePersistedConfigOptions: vi.fn(), + adapterStore: { + getAdapter: () => undefined, + setAdapter: vi.fn(), + removeAdapter: vi.fn(), + }, + trpc: { + agent: { + cancel: { mutate: cancelMutate }, + onSessionIdleKilled: { + subscribe: () => ({ unsubscribe: vi.fn() }), + }, + }, + }, + } as unknown as SessionServiceDeps; + + const service = new SessionService(deps); + return { service, sessions, removeSession, cancelMutate }; +} + +function connectParamsFor(taskId: string) { + return { + task: { id: taskId, title: taskId, description: taskId } as Task, + repoPath: "/repo", + }; +} + +describe("SessionService idle session eviction", () => { + it("evicts the least recently used idle sessions beyond the budget", async () => { + const idleCount = MAX_CONNECTED_SESSIONS; + const seeds = Array.from({ length: idleCount }, (_, i) => + makeSession(`idle-${i}`, i + 1), + ); + seeds.push(makeSession("active", 1000)); + const { service, removeSession } = createHarness(seeds); + + await service.connectToTask(connectParamsFor("active")); + + await vi.waitFor(() => { + expect(removeSession).toHaveBeenCalledTimes(2); + }); + expect(removeSession).toHaveBeenCalledWith("run-idle-0"); + expect(removeSession).toHaveBeenCalledWith("run-idle-1"); + }); + + it("never evicts mounted or busy sessions", async () => { + const idleCount = MAX_CONNECTED_SESSIONS; + const seeds = Array.from({ length: idleCount }, (_, i) => + makeSession(`idle-${i}`, i + 1, { + isPromptPending: i === 1, + }), + ); + seeds.push(makeSession("active", 1000)); + const { service, removeSession } = createHarness(seeds); + + const unregister = service.registerMountedTask("idle-0"); + + await service.connectToTask(connectParamsFor("active")); + + await vi.waitFor(() => { + expect(removeSession).toHaveBeenCalledTimes(2); + }); + expect(removeSession).not.toHaveBeenCalledWith("run-idle-0"); + expect(removeSession).not.toHaveBeenCalledWith("run-idle-1"); + expect(removeSession).toHaveBeenCalledWith("run-idle-2"); + expect(removeSession).toHaveBeenCalledWith("run-idle-3"); + unregister(); + }); + + it("evicts nothing at or under the budget", async () => { + const seeds = Array.from({ length: MAX_CONNECTED_SESSIONS - 2 }, (_, i) => + makeSession(`idle-${i}`, i + 1), + ); + seeds.push(makeSession("active", 1000)); + const { service, removeSession } = createHarness(seeds); + + await service.connectToTask(connectParamsFor("active")); + await new Promise((resolve) => setTimeout(resolve, 0)); + + expect(removeSession).not.toHaveBeenCalled(); + }); +}); diff --git a/packages/ui/src/features/sessions/hooks/useSessionConnection.ts b/packages/ui/src/features/sessions/hooks/useSessionConnection.ts index 6320cecb01..69ba979fff 100644 --- a/packages/ui/src/features/sessions/hooks/useSessionConnection.ts +++ b/packages/ui/src/features/sessions/hooks/useSessionConnection.ts @@ -71,6 +71,10 @@ export function useSessionConnection({ sessionEventCount, ]); + useEffect(() => { + return sessionService.registerMountedTask(task.id); + }, [task.id, sessionService]); + useEffect(() => { if (!taskRunId) return; return sessionService.startActivityHeartbeat(taskRunId); From e3241217791fcbd6bab3f21555f23031e3909106 Mon Sep 17 00:00:00 2001 From: Charles Vien Date: Wed, 1 Jul 2026 17:41:33 -0700 Subject: [PATCH 09/20] extract shared task terminal teardown helper --- .../ui/src/features/archive/useArchiveTask.ts | 7 +- .../features/terminal/TerminalManager.test.ts | 79 +++++++++++++++++++ .../features/terminal/destroyTaskTerminals.ts | 7 ++ 3 files changed, 88 insertions(+), 5 deletions(-) create mode 100644 packages/ui/src/features/terminal/destroyTaskTerminals.ts diff --git a/packages/ui/src/features/archive/useArchiveTask.ts b/packages/ui/src/features/archive/useArchiveTask.ts index b214ef8727..605c4aa576 100644 --- a/packages/ui/src/features/archive/useArchiveTask.ts +++ b/packages/ui/src/features/archive/useArchiveTask.ts @@ -19,7 +19,7 @@ import { useHostTRPC } from "@posthog/host-router/react"; import { useCommandCenterStore } from "@posthog/ui/features/command-center/commandCenterStore"; import { useFocusStore } from "@posthog/ui/features/focus/focusStore"; import { pinnedTasksApi } from "@posthog/ui/features/sidebar/taskMetaApi"; -import { terminalManager } from "@posthog/ui/features/terminal/TerminalManager"; +import { destroyTaskTerminals } from "@posthog/ui/features/terminal/destroyTaskTerminals"; import { type TerminalState, useTerminalStore, @@ -103,10 +103,7 @@ function makeOrchestrationDeps( ([key]) => key === taskId || key.startsWith(`${taskId}-`), ), ), - clearTerminalStates: (taskId) => { - terminalManager.destroyForTask(taskId); - useTerminalStore.getState().clearTerminalStatesForTask(taskId); - }, + clearTerminalStates: (taskId) => destroyTaskTerminals(taskId), restoreTerminalStates: (states) => { useTerminalStore.setState((s) => ({ terminalStates: { diff --git a/packages/ui/src/features/terminal/TerminalManager.test.ts b/packages/ui/src/features/terminal/TerminalManager.test.ts index 41ecc7429a..b250e4b89b 100644 --- a/packages/ui/src/features/terminal/TerminalManager.test.ts +++ b/packages/ui/src/features/terminal/TerminalManager.test.ts @@ -9,6 +9,7 @@ const mocks = vi.hoisted(() => { const openExternal = vi.fn(); const logInfo = vi.fn(); const logError = vi.fn(); + const logWarn = vi.fn(); class MockTerminal { cols = 80; @@ -56,6 +57,7 @@ const mocks = vi.hoisted(() => { openExternal, logInfo, logError, + logWarn, MockTerminal, terminalInstances, }; @@ -77,6 +79,7 @@ vi.mock("@posthog/ui/shell/logger", () => ({ scope: () => ({ info: mocks.logInfo, error: mocks.logError, + warn: mocks.logWarn, }), }, })); @@ -101,6 +104,13 @@ vi.mock("@xterm/addon-web-links", () => ({ WebLinksAddon: class {}, })); +vi.mock("@xterm/addon-webgl", () => ({ + WebglAddon: class { + onContextLoss = vi.fn(); + dispose = vi.fn(); + }, +})); + vi.mock("@xterm/xterm", () => ({ Terminal: mocks.MockTerminal, })); @@ -169,3 +179,72 @@ describe("TerminalManager shell recovery", () => { }); }); }); + +describe("TerminalManager.destroyForTask", () => { + beforeEach(() => { + mocks.check.mockReset().mockResolvedValue(true); + mocks.create.mockReset().mockResolvedValue(undefined); + mocks.write.mockReset().mockResolvedValue(undefined); + mocks.resize.mockReset().mockResolvedValue(undefined); + mocks.terminalInstances.length = 0; + vi.stubGlobal( + "ResizeObserver", + class { + observe() {} + unobserve() {} + disconnect() {} + }, + ); + }); + + afterEach(() => { + for (const id of terminalManager.getSessionsByPrefix("")) { + terminalManager.destroy(id); + } + vi.unstubAllGlobals(); + }); + + it("destroys the task's main and action terminals only", () => { + terminalManager.create({ + sessionId: "sess-a", + persistenceKey: "task-1", + taskId: "task-1", + }); + terminalManager.create({ + sessionId: "sess-b", + persistenceKey: "task-1-action", + taskId: "task-1", + }); + terminalManager.create({ + sessionId: "sess-c", + persistenceKey: "task-10", + taskId: "task-10", + }); + + terminalManager.destroyForTask("task-1"); + + expect(terminalManager.getSessionsByPrefix("sess-")).toEqual(["sess-c"]); + expect(mocks.terminalInstances[0].dispose).toHaveBeenCalled(); + expect(mocks.terminalInstances[1].dispose).toHaveBeenCalled(); + expect(mocks.terminalInstances[2].dispose).not.toHaveBeenCalled(); + }); + + it("removes the parked terminal element from the DOM on destroy", () => { + terminalManager.create({ + sessionId: "sess-parked", + persistenceKey: "task-2", + taskId: "task-2", + }); + const host = document.createElement("div"); + document.body.appendChild(host); + terminalManager.attach("sess-parked", host); + terminalManager.detach("sess-parked"); + + const parking = document.getElementById("terminal-parking"); + expect(parking?.childElementCount).toBe(1); + + terminalManager.destroy("sess-parked"); + expect(parking?.childElementCount).toBe(0); + host.remove(); + }); +}); diff --git a/packages/ui/src/features/terminal/destroyTaskTerminals.ts b/packages/ui/src/features/terminal/destroyTaskTerminals.ts new file mode 100644 index 0000000000..c739ece6bd --- /dev/null +++ b/packages/ui/src/features/terminal/destroyTaskTerminals.ts @@ -0,0 +1,7 @@ +import { terminalManager } from "./TerminalManager"; +import { useTerminalStore } from "./terminalStore"; + +export function destroyTaskTerminals(taskId: string): void { + terminalManager.destroyForTask(taskId); + useTerminalStore.getState().clearTerminalStatesForTask(taskId); +} From 8b9d330a57828b1b793c3bb2c812d894156f8f43 Mon Sep 17 00:00:00 2001 From: Charles Vien Date: Wed, 1 Jul 2026 17:41:34 -0700 Subject: [PATCH 10/20] destroy terminals only after suspend succeeds --- .../ui/src/features/suspension/useSuspendTask.test.tsx | 9 +++++---- packages/ui/src/features/suspension/useSuspendTask.ts | 7 ++----- 2 files changed, 7 insertions(+), 9 deletions(-) diff --git a/packages/ui/src/features/suspension/useSuspendTask.test.tsx b/packages/ui/src/features/suspension/useSuspendTask.test.tsx index 248a2cf69a..2524a0699b 100644 --- a/packages/ui/src/features/suspension/useSuspendTask.test.tsx +++ b/packages/ui/src/features/suspension/useSuspendTask.test.tsx @@ -4,6 +4,7 @@ import type { ReactNode } from "react"; import { beforeEach, describe, expect, it, vi } from "vitest"; const suspendFn = vi.hoisted(() => vi.fn().mockResolvedValue(undefined)); +const destroyTaskTerminals = vi.hoisted(() => vi.fn()); const workspaceClient = vi.hoisted(() => ({ getAll: vi.fn().mockResolvedValue({}), })); @@ -32,10 +33,8 @@ vi.mock("@posthog/host-router/react", () => ({ vi.mock("@posthog/ui/features/focus/focusStore", () => ({ useFocusStore: { getState: () => ({ session: null, disableFocus: vi.fn() }) }, })); -vi.mock("@posthog/ui/features/terminal/terminalStore", () => ({ - useTerminalStore: { - getState: () => ({ clearTerminalStatesForTask: vi.fn() }), - }, +vi.mock("@posthog/ui/features/terminal/destroyTaskTerminals", () => ({ + destroyTaskTerminals, })); vi.mock("@posthog/ui/shell/logger", () => ({ logger: { scope: () => ({ info: vi.fn(), error: vi.fn() }) }, @@ -66,6 +65,7 @@ describe("useSuspendTask", () => { taskId: "t1", reason: "manual", }); + expect(destroyTaskTerminals).toHaveBeenCalledWith("t1"); }); it("rolls back the optimistic suspended set when suspend fails", async () => { @@ -86,5 +86,6 @@ describe("useSuspendTask", () => { ); seen.push(queryClient.getQueryData(SUSPENDED_TASK_IDS_KEY)); expect(seen[0] ?? []).not.toContain("t1"); + expect(destroyTaskTerminals).not.toHaveBeenCalled(); }); }); diff --git a/packages/ui/src/features/suspension/useSuspendTask.ts b/packages/ui/src/features/suspension/useSuspendTask.ts index c9e9b0c7b8..7bb88f385c 100644 --- a/packages/ui/src/features/suspension/useSuspendTask.ts +++ b/packages/ui/src/features/suspension/useSuspendTask.ts @@ -1,7 +1,6 @@ import { useHostTRPC, useHostTRPCClient } from "@posthog/host-router/react"; import { useFocusStore } from "@posthog/ui/features/focus/focusStore"; -import { terminalManager } from "@posthog/ui/features/terminal/TerminalManager"; -import { useTerminalStore } from "@posthog/ui/features/terminal/terminalStore"; +import { destroyTaskTerminals } from "@posthog/ui/features/terminal/destroyTaskTerminals"; import { logger } from "@posthog/ui/shell/logger"; import { useMutation, useQueryClient } from "@tanstack/react-query"; import { WORKSPACE_QUERY_KEY } from "../workspace/identifiers"; @@ -30,9 +29,6 @@ export function useSuspendTask() { const workspaces = await hostClient.workspace.getAll.query(); const workspace = workspaces[taskId] ?? null; - terminalManager.destroyForTask(taskId); - useTerminalStore.getState().clearTerminalStatesForTask(taskId); - queryClient.setQueryData(suspendedTaskIdsKey, (old) => old ? [...old, taskId] : [taskId], ); @@ -48,6 +44,7 @@ export function useSuspendTask() { try { await suspendMutation.mutateAsync({ taskId, reason }); + destroyTaskTerminals(taskId); queryClient.invalidateQueries({ queryKey: suspensionPathKey }); queryClient.invalidateQueries({ queryKey: suspendedTaskIdsKey }); queryClient.invalidateQueries({ queryKey: WORKSPACE_QUERY_KEY }); From 710a570c07cc1822b7a1e801836dc879bb7aa75b Mon Sep 17 00:00:00 2001 From: Charles Vien Date: Wed, 1 Jul 2026 17:41:36 -0700 Subject: [PATCH 11/20] harden delete cleanup and drop resolveService --- .../tasks/useTaskCrudMutations.test.tsx | 63 ++++++++++++++++++- .../features/tasks/useTaskCrudMutations.ts | 25 +++++--- 2 files changed, 77 insertions(+), 11 deletions(-) diff --git a/packages/ui/src/features/tasks/useTaskCrudMutations.test.tsx b/packages/ui/src/features/tasks/useTaskCrudMutations.test.tsx index c63b4c9c68..7270003412 100644 --- a/packages/ui/src/features/tasks/useTaskCrudMutations.test.tsx +++ b/packages/ui/src/features/tasks/useTaskCrudMutations.test.tsx @@ -21,15 +21,25 @@ const deletionService = vi.hoisted(() => ({ confirmAndDelete, })); +const destroyTaskTerminals = vi.hoisted(() => vi.fn()); + vi.mock("@posthog/ui/hooks/useAuthenticatedMutation", () => ({ useAuthenticatedMutation: () => ({ mutateAsync, isPending: false }), })); vi.mock("@posthog/di/react", () => ({ useService: () => deletionService, })); +vi.mock("@posthog/ui/features/terminal/destroyTaskTerminals", () => ({ + destroyTaskTerminals, +})); +import type { SessionService } from "@posthog/core/sessions/sessionService"; import { taskKeys } from "./taskKeys"; -import { useCreateTask, useDeleteTask } from "./useTaskCrudMutations"; +import { + releaseDeletedTaskResources, + useCreateTask, + useDeleteTask, +} from "./useTaskCrudMutations"; function wrapper({ children }: { children: ReactNode }) { const queryClient = new QueryClient(); @@ -88,6 +98,57 @@ describe("useDeleteTask.deleteWithConfirm", () => { }); }); +describe("releaseDeletedTaskResources", () => { + beforeEach(() => { + vi.clearAllMocks(); + }); + + it("disconnects the session and destroys the task's terminals", async () => { + const sessionService = { + disconnectFromTask: vi.fn().mockResolvedValue(undefined), + }; + + await releaseDeletedTaskResources( + "t1", + sessionService as unknown as SessionService, + ); + + expect(sessionService.disconnectFromTask).toHaveBeenCalledWith("t1"); + expect(destroyTaskTerminals).toHaveBeenCalledWith("t1"); + }); + + it("still destroys terminals when the disconnect fails, without throwing", async () => { + const sessionService = { + disconnectFromTask: vi.fn().mockRejectedValue(new Error("gone")), + }; + + await expect( + releaseDeletedTaskResources( + "t1", + sessionService as unknown as SessionService, + ), + ).resolves.toBeUndefined(); + + expect(destroyTaskTerminals).toHaveBeenCalledWith("t1"); + }); + + it("never rejects when terminal teardown fails", async () => { + destroyTaskTerminals.mockImplementationOnce(() => { + throw new Error("boom"); + }); + const sessionService = { + disconnectFromTask: vi.fn().mockResolvedValue(undefined), + }; + + await expect( + releaseDeletedTaskResources( + "t1", + sessionService as unknown as SessionService, + ), + ).resolves.toBeUndefined(); + }); +}); + describe("useCreateTask.invalidateTasks", () => { beforeEach(() => { vi.clearAllMocks(); diff --git a/packages/ui/src/features/tasks/useTaskCrudMutations.ts b/packages/ui/src/features/tasks/useTaskCrudMutations.ts index e3f13044b1..628dde456f 100644 --- a/packages/ui/src/features/tasks/useTaskCrudMutations.ts +++ b/packages/ui/src/features/tasks/useTaskCrudMutations.ts @@ -10,11 +10,9 @@ import { TASK_DELETION_SERVICE, type TaskDeletionService, } from "@posthog/core/tasks/taskDeletionService"; -import { resolveService } from "@posthog/di/container"; import { useService } from "@posthog/di/react"; import type { Task } from "@posthog/shared/domain-types"; -import { terminalManager } from "@posthog/ui/features/terminal/TerminalManager"; -import { useTerminalStore } from "@posthog/ui/features/terminal/terminalStore"; +import { destroyTaskTerminals } from "@posthog/ui/features/terminal/destroyTaskTerminals"; import { useAuthenticatedMutation } from "@posthog/ui/hooks/useAuthenticatedMutation"; import { logger } from "@posthog/ui/shell/logger"; import { useQueryClient } from "@tanstack/react-query"; @@ -23,16 +21,22 @@ import { taskKeys } from "./taskKeys"; const log = logger.scope("tasks"); -async function releaseDeletedTaskResources(taskId: string): Promise { +// Never throws: the task is already deleted server-side, so a cleanup failure +// must not reject the mutation and roll back the optimistic list removal. +export async function releaseDeletedTaskResources( + taskId: string, + sessionService: SessionService, +): Promise { try { - await resolveService(SESSION_SERVICE).disconnectFromTask( - taskId, - ); + await sessionService.disconnectFromTask(taskId); } catch (error) { log.error("Failed to disconnect session for deleted task", error); } - terminalManager.destroyForTask(taskId); - useTerminalStore.getState().clearTerminalStatesForTask(taskId); + try { + destroyTaskTerminals(taskId); + } catch (error) { + log.error("Failed to release terminals for deleted task", error); + } } export function useCreateTask() { @@ -97,11 +101,12 @@ export function useDeleteTask() { const deletionService = useService( TASK_DELETION_SERVICE, ); + const sessionService = useService(SESSION_SERVICE); const mutation = useAuthenticatedMutation( async (client, taskId: string) => { const result = await deletionService.deleteTask(client, taskId); - await releaseDeletedTaskResources(taskId); + await releaseDeletedTaskResources(taskId, sessionService); return result; }, { From b78ab0b1baa18070f76824de973bf7115838a693 Mon Sep 17 00:00:00 2001 From: Charles Vien Date: Wed, 1 Jul 2026 17:41:37 -0700 Subject: [PATCH 12/20] note sequential probing and test short-circuit --- .../src/services/connectivity/service.test.ts | 14 +++++++++++++- .../src/services/connectivity/service.ts | 2 ++ 2 files changed, 15 insertions(+), 1 deletion(-) diff --git a/packages/workspace-server/src/services/connectivity/service.test.ts b/packages/workspace-server/src/services/connectivity/service.test.ts index 44be28df52..72fd674929 100644 --- a/packages/workspace-server/src/services/connectivity/service.test.ts +++ b/packages/workspace-server/src/services/connectivity/service.test.ts @@ -120,7 +120,7 @@ describe("ConnectivityService", () => { service.on(ConnectivityEvent.StatusChange, handler); mockFetch.mockRejectedValue(new Error("offline")); - await vi.advanceTimersByTimeAsync(3000); // one failed poll only + await vi.advanceTimersByTimeAsync(30_000); // one failed poll only expect(handler).not.toHaveBeenCalled(); }); @@ -201,6 +201,18 @@ describe("ConnectivityService", () => { expect(result).toEqual({ isOnline: true }); }); + it("short-circuits after the first reachable host", async () => { + mockFetch.mockResolvedValue(ok(204)); + service = new ConnectivityService(); + await vi.advanceTimersByTimeAsync(0); + + expect(mockFetch).toHaveBeenCalledTimes(1); + expect(mockFetch).toHaveBeenCalledWith( + "https://www.google.com/generate_204", + expect.objectContaining({ method: "HEAD" }), + ); + }); + it("goes offline only when every host fails", async () => { mockFetch.mockRejectedValue(new Error("blocked")); diff --git a/packages/workspace-server/src/services/connectivity/service.ts b/packages/workspace-server/src/services/connectivity/service.ts index 114c25fbe2..7f4434ef77 100644 --- a/packages/workspace-server/src/services/connectivity/service.ts +++ b/packages/workspace-server/src/services/connectivity/service.ts @@ -75,6 +75,8 @@ export class ConnectivityService extends TypedEventEmitter { } private async verifyWithHttp(): Promise { + // Sequential on purpose: one request per check in the common case, at the + // cost of up to CHECK_TIMEOUT_MS extra latency when the first host is blocked. for (const url of CHECK_URLS) { try { await this.probe(url); From ddc646b496666b154d2b123f77a6f3e369d9696d Mon Sep 17 00:00:00 2001 From: Charles Vien Date: Wed, 1 Jul 2026 18:10:53 -0700 Subject: [PATCH 13/20] align event cap with transcript evict and restore --- .../core/src/sessions/sessionStore.test.ts | 37 +++++++++++++++++++ packages/core/src/sessions/sessionStore.ts | 15 ++++++-- 2 files changed, 49 insertions(+), 3 deletions(-) diff --git a/packages/core/src/sessions/sessionStore.test.ts b/packages/core/src/sessions/sessionStore.test.ts index 5d565b93eb..a1e84b7914 100644 --- a/packages/core/src/sessions/sessionStore.test.ts +++ b/packages/core/src/sessions/sessionStore.test.ts @@ -120,4 +120,41 @@ describe("sessionStore event cap", () => { expect(session.events).toHaveLength(10); expect(session.trimmedEventCount ?? 0).toBe(0); }); + + it("resets the trim offset when updateSession replaces the stream", () => { + sessionStoreSetters.appendEvents("run-1", events(MAX_SESSION_EVENTS + 40)); + expect(sessionStore.getState().sessions["run-1"].trimmedEventCount).toBe( + 40, + ); + + sessionStoreSetters.updateSession("run-1", { events: events(10) }); + + const session = sessionStore.getState().sessions["run-1"]; + expect(session.events).toHaveLength(10); + expect(session.trimmedEventCount).toBe(0); + }); + + it("caps rehydrated histories in restoreEvents and restarts the offset", () => { + sessionStoreSetters.appendEvents("run-1", events(MAX_SESSION_EVENTS + 40)); + + sessionStoreSetters.restoreEvents( + "run-1", + events(MAX_SESSION_EVENTS + 7), + MAX_SESSION_EVENTS + 7, + ); + + const session = sessionStore.getState().sessions["run-1"]; + expect(session.events).toHaveLength(MAX_SESSION_EVENTS); + expect(session.trimmedEventCount).toBe(7); + }); + + it("clears the trim offset when a transcript is evicted", () => { + sessionStoreSetters.appendEvents("run-1", events(MAX_SESSION_EVENTS + 5)); + + sessionStoreSetters.evictEvents("run-1"); + + const session = sessionStore.getState().sessions["run-1"]; + expect(session.events).toHaveLength(0); + expect(session.trimmedEventCount).toBe(0); + }); }); diff --git a/packages/core/src/sessions/sessionStore.ts b/packages/core/src/sessions/sessionStore.ts index 8054740500..8f59ac1c13 100644 --- a/packages/core/src/sessions/sessionStore.ts +++ b/packages/core/src/sessions/sessionStore.ts @@ -69,10 +69,16 @@ export const sessionStoreSetters = { updateSession: (taskRunId: string, updates: Partial) => { sessionStore.setState((state) => { - if (state.sessions[taskRunId]) { - Object.assign(state.sessions[taskRunId], updates); + const session = state.sessions[taskRunId]; + if (session) { + Object.assign(session, updates); if (updates.events) { - trimSessionEvents(state.sessions[taskRunId]); + // A wholesale replacement starts a new stream snapshot, so a stale + // head-trim offset from the previous array must not carry over. + if (updates.trimmedEventCount === undefined) { + session.trimmedEventCount = 0; + } + trimSessionEvents(session); } } }); @@ -109,6 +115,7 @@ export const sessionStoreSetters = { if (session && session.events.length > 0) { session.events = []; session.processedLineCount = 0; + session.trimmedEventCount = 0; } }); }, @@ -128,6 +135,8 @@ export const sessionStoreSetters = { for (const event of events) Object.freeze(event); session.events = events; session.processedLineCount = lineCount; + session.trimmedEventCount = 0; + trimSessionEvents(session); } }); }, From 191d79fa760da4aeff382c7819dbb84a400e6f29 Mon Sep 17 00:00:00 2001 From: Charles Vien Date: Wed, 1 Jul 2026 19:19:13 -0700 Subject: [PATCH 14/20] drop the session event cap --- .../src/sessions/cloudRunIdleTracker.test.ts | 36 ---- .../core/src/sessions/cloudRunIdleTracker.ts | 30 +--- .../core/src/sessions/sessionStore.test.ts | 160 ------------------ packages/core/src/sessions/sessionStore.ts | 29 +--- packages/shared/src/sessions.ts | 1 - 5 files changed, 9 insertions(+), 247 deletions(-) delete mode 100644 packages/core/src/sessions/sessionStore.test.ts diff --git a/packages/core/src/sessions/cloudRunIdleTracker.test.ts b/packages/core/src/sessions/cloudRunIdleTracker.test.ts index abeb4ee135..49c7ae2c13 100644 --- a/packages/core/src/sessions/cloudRunIdleTracker.test.ts +++ b/packages/core/src/sessions/cloudRunIdleTracker.test.ts @@ -130,39 +130,3 @@ describe("CloudRunIdleTracker mark/capture/restore", () => { expect(tracker.evaluateIdle(session("r1", [])).idle).toBe(false); }); }); - -describe("CloudRunIdleTracker with trimmed event history", () => { - it("keeps incremental cursors stable when the event head is trimmed", () => { - const tracker = new CloudRunIdleTracker(); - const s = session("r1", [runStarted("r1"), promptRequest()]); - expect(tracker.evaluateIdle(s).idle).toBe(false); - - const trimmed = session("r1", [turnComplete()]); - trimmed.trimmedEventCount = 2; - expect(tracker.evaluateIdle(trimmed).idle).toBe(true); - }); - - it("does not skip events appended in the same batch as a trim", () => { - const tracker = new CloudRunIdleTracker(); - const s = session("r1", [runStarted("r1")]); - expect(tracker.evaluateIdle(s).idle).toBe(false); - - const trimmed = session("r1", [promptRequest(), turnComplete()]); - trimmed.trimmedEventCount = 1; - expect(tracker.evaluateIdle(trimmed).idle).toBe(true); - }); - - it("restoreAfterFailedSend translates snapshot offsets across trims", () => { - const tracker = new CloudRunIdleTracker(); - const before = session("r1", [runStarted("r1")], "r1"); - tracker.markIdle(before); - const snapshot = tracker.capture(before); - - const after = session("r1", [], "r1"); - after.trimmedEventCount = 1; - tracker.markBusy(after); - expect(tracker.restoreAfterFailedSend(snapshot, after)).toEqual({ - agentIdleForRunId: "r1", - }); - }); -}); diff --git a/packages/core/src/sessions/cloudRunIdleTracker.ts b/packages/core/src/sessions/cloudRunIdleTracker.ts index 8afac6c159..24ef58e367 100644 --- a/packages/core/src/sessions/cloudRunIdleTracker.ts +++ b/packages/core/src/sessions/cloudRunIdleTracker.ts @@ -24,14 +24,6 @@ export interface CloudRunIdleScanResult { shouldCacheToStore: boolean; } -function streamLength(session: AgentSession): number { - return (session.trimmedEventCount ?? 0) + session.events.length; -} - -function toArrayIndex(streamOffset: number, session: AgentSession): number { - return Math.max(0, streamOffset - (session.trimmedEventCount ?? 0)); -} - /** * Tracks idleness for cloud runs incrementally so repeated `in_progress` * updates don't re-scan the full event list each time. @@ -54,7 +46,7 @@ export class CloudRunIdleTracker { */ markBusy(session: AgentSession): void { this.scanStates.set(session.taskRunId, { - nextEventIndex: streamLength(session), + nextEventIndex: session.events.length, seenCurrentRunStart: true, idle: false, }); @@ -62,7 +54,7 @@ export class CloudRunIdleTracker { markIdle(session: AgentSession): void { this.scanStates.set(session.taskRunId, { - nextEventIndex: streamLength(session), + nextEventIndex: session.events.length, seenCurrentRunStart: true, idle: true, }); @@ -72,7 +64,7 @@ export class CloudRunIdleTracker { const scanState = this.scanStates.get(session.taskRunId); return { taskRunId: session.taskRunId, - eventCount: streamLength(session), + eventCount: session.events.length, agentIdleForRunId: session.agentIdleForRunId, scanState: scanState ? { ...scanState } : undefined, }; @@ -82,11 +74,7 @@ export class CloudRunIdleTracker { snapshot: CloudRunIdleEvidenceSnapshot, session: AgentSession, ): CloudRunIdleRestoreResult | undefined { - for ( - let i = toArrayIndex(snapshot.eventCount, session); - i < session.events.length; - i += 1 - ) { + for (let i = snapshot.eventCount; i < session.events.length; i += 1) { const acpMsg = session.events[i]; if ( acpMsg && @@ -126,7 +114,7 @@ export class CloudRunIdleTracker { } let scanState = this.scanStates.get(session.taskRunId); - if (!scanState || scanState.nextEventIndex > streamLength(session)) { + if (!scanState || scanState.nextEventIndex > session.events.length) { scanState = { nextEventIndex: 0, seenCurrentRunStart: false, @@ -134,11 +122,7 @@ export class CloudRunIdleTracker { }; } - for ( - let i = toArrayIndex(scanState.nextEventIndex, session); - i < session.events.length; - i += 1 - ) { + for (let i = scanState.nextEventIndex; i < session.events.length; i += 1) { const acpMsg = session.events[i]; if (!acpMsg) continue; const msg = acpMsg.message; @@ -168,7 +152,7 @@ export class CloudRunIdleTracker { } } - scanState.nextEventIndex = streamLength(session); + scanState.nextEventIndex = session.events.length; this.scanStates.set(session.taskRunId, scanState); return { idle: scanState.idle, shouldCacheToStore: scanState.idle }; diff --git a/packages/core/src/sessions/sessionStore.test.ts b/packages/core/src/sessions/sessionStore.test.ts deleted file mode 100644 index a1e84b7914..0000000000 --- a/packages/core/src/sessions/sessionStore.test.ts +++ /dev/null @@ -1,160 +0,0 @@ -import type { AcpMessage, AgentSession } from "@posthog/shared"; -import { beforeEach, describe, expect, it } from "vitest"; -import { - MAX_SESSION_EVENTS, - sessionStore, - sessionStoreSetters, -} from "./sessionStore"; - -function event(n: number): AcpMessage { - return { - type: "acp_message", - ts: n, - message: { - jsonrpc: "2.0", - method: "session/update", - params: { sessionUpdate: "agent_message_chunk", n }, - }, - } as AcpMessage; -} - -function makeSession(taskRunId: string, taskId: string): AgentSession { - return { - taskRunId, - taskId, - events: [], - optimisticItems: [], - messageQueue: [], - pendingPermissions: new Map(), - } as unknown as AgentSession; -} - -function events(count: number, offset = 0): AcpMessage[] { - return Array.from({ length: count }, (_, i) => event(offset + i)); -} - -describe("sessionStore event cap", () => { - beforeEach(() => { - sessionStoreSetters.clearAll(); - sessionStoreSetters.setSession(makeSession("run-1", "task-1")); - }); - - it("keeps events unbounded-free: appends under the cap verbatim", () => { - sessionStoreSetters.appendEvents("run-1", events(10)); - - const session = sessionStore.getState().sessions["run-1"]; - expect(session.events).toHaveLength(10); - expect(session.trimmedEventCount ?? 0).toBe(0); - }); - - it("drops the oldest events beyond MAX_SESSION_EVENTS and records the trim", () => { - sessionStoreSetters.appendEvents("run-1", events(MAX_SESSION_EVENTS + 100)); - - const session = sessionStore.getState().sessions["run-1"]; - expect(session.events).toHaveLength(MAX_SESSION_EVENTS); - expect(session.trimmedEventCount).toBe(100); - expect( - (session.events[0].message as { params: { n: number } }).params.n, - ).toBe(100); - }); - - it("accumulates trimmedEventCount across appends", () => { - sessionStoreSetters.appendEvents("run-1", events(MAX_SESSION_EVENTS)); - sessionStoreSetters.appendEvents("run-1", events(50, MAX_SESSION_EVENTS)); - sessionStoreSetters.appendEvents( - "run-1", - events(25, MAX_SESSION_EVENTS + 50), - ); - - const session = sessionStore.getState().sessions["run-1"]; - expect(session.events).toHaveLength(MAX_SESSION_EVENTS); - expect(session.trimmedEventCount).toBe(75); - expect( - (session.events.at(-1)?.message as { params: { n: number } }).params.n, - ).toBe(MAX_SESSION_EVENTS + 74); - }); - - it("trims in replaceOptimisticWithEvent too", () => { - sessionStoreSetters.appendEvents("run-1", events(MAX_SESSION_EVENTS)); - sessionStoreSetters.appendOptimisticItem("run-1", { - type: "user_message", - content: "hi", - } as never); - - sessionStoreSetters.replaceOptimisticWithEvent( - "run-1", - event(MAX_SESSION_EVENTS), - ); - - const session = sessionStore.getState().sessions["run-1"]; - expect(session.events).toHaveLength(MAX_SESSION_EVENTS); - expect(session.trimmedEventCount).toBe(1); - expect(session.optimisticItems).toHaveLength(0); - }); - - it("trims oversized histories installed via setSession", () => { - const hydrated = makeSession("run-2", "task-2"); - hydrated.events = events(MAX_SESSION_EVENTS + 200); - sessionStoreSetters.setSession(hydrated); - - const session = sessionStore.getState().sessions["run-2"]; - expect(session.events).toHaveLength(MAX_SESSION_EVENTS); - expect(session.trimmedEventCount).toBe(200); - }); - - it("trims oversized histories installed via updateSession", () => { - sessionStoreSetters.updateSession("run-1", { - events: events(MAX_SESSION_EVENTS + 30), - }); - - const session = sessionStore.getState().sessions["run-1"]; - expect(session.events).toHaveLength(MAX_SESSION_EVENTS); - expect(session.trimmedEventCount).toBe(30); - }); - - it("leaves untouched events alone in updateSession without events", () => { - sessionStoreSetters.appendEvents("run-1", events(10)); - sessionStoreSetters.updateSession("run-1", { taskTitle: "renamed" }); - - const session = sessionStore.getState().sessions["run-1"]; - expect(session.events).toHaveLength(10); - expect(session.trimmedEventCount ?? 0).toBe(0); - }); - - it("resets the trim offset when updateSession replaces the stream", () => { - sessionStoreSetters.appendEvents("run-1", events(MAX_SESSION_EVENTS + 40)); - expect(sessionStore.getState().sessions["run-1"].trimmedEventCount).toBe( - 40, - ); - - sessionStoreSetters.updateSession("run-1", { events: events(10) }); - - const session = sessionStore.getState().sessions["run-1"]; - expect(session.events).toHaveLength(10); - expect(session.trimmedEventCount).toBe(0); - }); - - it("caps rehydrated histories in restoreEvents and restarts the offset", () => { - sessionStoreSetters.appendEvents("run-1", events(MAX_SESSION_EVENTS + 40)); - - sessionStoreSetters.restoreEvents( - "run-1", - events(MAX_SESSION_EVENTS + 7), - MAX_SESSION_EVENTS + 7, - ); - - const session = sessionStore.getState().sessions["run-1"]; - expect(session.events).toHaveLength(MAX_SESSION_EVENTS); - expect(session.trimmedEventCount).toBe(7); - }); - - it("clears the trim offset when a transcript is evicted", () => { - sessionStoreSetters.appendEvents("run-1", events(MAX_SESSION_EVENTS + 5)); - - sessionStoreSetters.evictEvents("run-1"); - - const session = sessionStore.getState().sessions["run-1"]; - expect(session.events).toHaveLength(0); - expect(session.trimmedEventCount).toBe(0); - }); -}); diff --git a/packages/core/src/sessions/sessionStore.ts b/packages/core/src/sessions/sessionStore.ts index 8f59ac1c13..0e86a78c6a 100644 --- a/packages/core/src/sessions/sessionStore.ts +++ b/packages/core/src/sessions/sessionStore.ts @@ -32,16 +32,6 @@ export const sessionStore = createStore()( })), ); -export const MAX_SESSION_EVENTS = 5000; - -function trimSessionEvents(session: AgentSession) { - const excess = session.events.length - MAX_SESSION_EVENTS; - if (excess > 0) { - session.events.splice(0, excess); - session.trimmedEventCount = (session.trimmedEventCount ?? 0) + excess; - } -} - export const sessionStoreSetters = { setSession: (session: AgentSession) => { sessionStore.setState((state) => { @@ -53,7 +43,6 @@ export const sessionStoreSetters = { state.sessions[session.taskRunId] = session; state.taskIdIndex[session.taskId] = session.taskRunId; - trimSessionEvents(state.sessions[session.taskRunId]); }); }, @@ -69,17 +58,8 @@ export const sessionStoreSetters = { updateSession: (taskRunId: string, updates: Partial) => { sessionStore.setState((state) => { - const session = state.sessions[taskRunId]; - if (session) { - Object.assign(session, updates); - if (updates.events) { - // A wholesale replacement starts a new stream snapshot, so a stale - // head-trim offset from the previous array must not carry over. - if (updates.trimmedEventCount === undefined) { - session.trimmedEventCount = 0; - } - trimSessionEvents(session); - } + if (state.sessions[taskRunId]) { + Object.assign(state.sessions[taskRunId], updates); } }); }, @@ -96,7 +76,6 @@ export const sessionStoreSetters = { // immer autofreeze, so this is the only freeze. for (const event of events) Object.freeze(event); session.events.push(...events); - trimSessionEvents(session); if (newLineCount !== undefined) { session.processedLineCount = newLineCount; } @@ -115,7 +94,6 @@ export const sessionStoreSetters = { if (session && session.events.length > 0) { session.events = []; session.processedLineCount = 0; - session.trimmedEventCount = 0; } }); }, @@ -135,8 +113,6 @@ export const sessionStoreSetters = { for (const event of events) Object.freeze(event); session.events = events; session.processedLineCount = lineCount; - session.trimmedEventCount = 0; - trimSessionEvents(session); } }); }, @@ -325,7 +301,6 @@ export const sessionStoreSetters = { const session = state.sessions[taskRunId]; if (session) { session.events.push(Object.freeze(event)); - trimSessionEvents(session); session.optimisticItems = []; } }); diff --git a/packages/shared/src/sessions.ts b/packages/shared/src/sessions.ts index 41728de954..0724dddfac 100644 --- a/packages/shared/src/sessions.ts +++ b/packages/shared/src/sessions.ts @@ -62,7 +62,6 @@ export interface AgentSession { currentPromptId?: number | null; logUrl?: string; processedLineCount?: number; - trimmedEventCount?: number; framework?: "claude"; adapter?: Adapter; configOptions?: SessionConfigOption[]; From 14148f66a6bf29cde2593420885ef2d387f68801 Mon Sep 17 00:00:00 2001 From: Charles Vien Date: Wed, 1 Jul 2026 19:20:05 -0700 Subject: [PATCH 15/20] preserve resume state when evicting idle sessions --- packages/core/src/sessions/sessionService.ts | 17 +++++++++++++---- 1 file changed, 13 insertions(+), 4 deletions(-) diff --git a/packages/core/src/sessions/sessionService.ts b/packages/core/src/sessions/sessionService.ts index 30791b4be4..431798630d 100644 --- a/packages/core/src/sessions/sessionService.ts +++ b/packages/core/src/sessions/sessionService.ts @@ -1034,7 +1034,10 @@ export class SessionService { } } - private async teardownSession(taskRunId: string): Promise { + private async teardownSession( + taskRunId: string, + opts?: { preserveResumeState?: boolean }, + ): Promise { const session = this.getSessionByRunId(taskRunId); try { @@ -1060,8 +1063,12 @@ export class SessionService { this.localRecoveryAttempts.delete(session.taskId); this.sessionLastUsedAt.delete(session.taskId); } - this.d.adapterStore.removeAdapter(taskRunId); - this.d.removePersistedConfigOptions(taskRunId); + if (!opts?.preserveResumeState) { + // Reconnect restores the model and permission mode from these; only a + // permanent disconnect (archive, delete, fresh session) may drop them. + this.d.adapterStore.removeAdapter(taskRunId); + this.d.removePersistedConfigOptions(taskRunId); + } } /** @@ -1397,7 +1404,9 @@ export class SessionService { }); this.sessionLastUsedAt.delete(session.taskId); try { - await this.teardownSession(session.taskRunId); + await this.teardownSession(session.taskRunId, { + preserveResumeState: true, + }); } catch (error) { this.d.log.error("Failed to evict idle session", { taskId: session.taskId, From bf0197fc2a89a38ab08815d5ad1462561aa12423 Mon Sep 17 00:00:00 2001 From: Charles Vien Date: Wed, 1 Jul 2026 19:20:55 -0700 Subject: [PATCH 16/20] evict idle sessions on cloud reconcile --- packages/core/src/sessions/sessionService.ts | 4 ++++ 1 file changed, 4 insertions(+) diff --git a/packages/core/src/sessions/sessionService.ts b/packages/core/src/sessions/sessionService.ts index 431798630d..36b18a1bd6 100644 --- a/packages/core/src/sessions/sessionService.ts +++ b/packages/core/src/sessions/sessionService.ts @@ -4344,6 +4344,10 @@ export class SessionService { } = params; if (isCloud) { + // Local connects bound the session budget inside connectToTask; cloud + // watches would otherwise never trigger eviction. + this.sessionLastUsedAt.set(task.id, Date.now()); + void this.evictIdleSessions(task.id); return this.reconcileCloudConnection( task, cloudAuth, From 9a594992d318021ecb3d52a510a96a8b23ac5c7f Mon Sep 17 00:00:00 2001 From: Charles Vien Date: Wed, 1 Jul 2026 19:20:57 -0700 Subject: [PATCH 17/20] treat compacting and handoff as busy for eviction --- packages/core/src/sessions/sessionEviction.ts | 4 ++++ 1 file changed, 4 insertions(+) diff --git a/packages/core/src/sessions/sessionEviction.ts b/packages/core/src/sessions/sessionEviction.ts index d56012d875..71b0879d12 100644 --- a/packages/core/src/sessions/sessionEviction.ts +++ b/packages/core/src/sessions/sessionEviction.ts @@ -7,6 +7,8 @@ export const MAX_CONNECTED_SESSIONS = 12; export function isSessionIdle(session: AgentSession): boolean { if (session.status === "connecting") return false; if (session.isPromptPending) return false; + if (session.isCompacting) return false; + if (session.handoffInProgress) return false; if (session.pendingPermissions.size > 0) return false; if (session.messageQueue.length > 0) return false; if (session.isCloud) return isTerminalStatus(session.cloudStatus); @@ -23,6 +25,8 @@ export function selectSessionsToEvict(params: { const { sessions, activeTaskId, protectedTaskIds, lastUsedAt } = params; const maxSessions = params.maxSessions ?? MAX_CONNECTED_SESSIONS; + // Reserves a slot for the incoming session even when a resume replaces an + // existing one; deliberately over-evicts by one in that case. const excess = sessions.length - (maxSessions - 1); if (excess <= 0) return []; From e1a2a9e0629d1137774496db8d862ff6711920fd Mon Sep 17 00:00:00 2001 From: Charles Vien Date: Wed, 1 Jul 2026 19:21:14 -0700 Subject: [PATCH 18/20] match task terminals by tagged task id --- .../features/terminal/TerminalManager.test.ts | 19 ++++++++++++++++++- .../src/features/terminal/TerminalManager.ts | 9 ++++++++- 2 files changed, 26 insertions(+), 2 deletions(-) diff --git a/packages/ui/src/features/terminal/TerminalManager.test.ts b/packages/ui/src/features/terminal/TerminalManager.test.ts index b250e4b89b..c647478944 100644 --- a/packages/ui/src/features/terminal/TerminalManager.test.ts +++ b/packages/ui/src/features/terminal/TerminalManager.test.ts @@ -212,7 +212,9 @@ describe("TerminalManager.destroyForTask", () => { }); terminalManager.create({ sessionId: "sess-b", - persistenceKey: "task-1-action", + // Production action-terminal key shape: the taskId sits mid-key, so + // only the tagged instance.taskId can match it. + persistenceKey: "action-setup-task-1-1700000000000-0", taskId: "task-1", }); terminalManager.create({ @@ -229,6 +231,21 @@ describe("TerminalManager.destroyForTask", () => { expect(mocks.terminalInstances[2].dispose).not.toHaveBeenCalled(); }); + it("falls back to the persistence key when an instance has no taskId", () => { + terminalManager.create({ + sessionId: "sess-d", + persistenceKey: "task-3-shell", + }); + terminalManager.create({ + sessionId: "sess-e", + persistenceKey: "task-30-shell", + }); + + terminalManager.destroyForTask("task-3"); + + expect(terminalManager.getSessionsByPrefix("sess-")).toEqual(["sess-e"]); + }); + it("removes the parked terminal element from the DOM on destroy", () => { terminalManager.create({ sessionId: "sess-parked", diff --git a/packages/ui/src/features/terminal/TerminalManager.ts b/packages/ui/src/features/terminal/TerminalManager.ts index f91ed19556..c6fd0cb66a 100644 --- a/packages/ui/src/features/terminal/TerminalManager.ts +++ b/packages/ui/src/features/terminal/TerminalManager.ts @@ -573,8 +573,15 @@ class TerminalManagerImpl { destroyForTask(taskId: string): void { for (const [sessionId, instance] of this.instances) { + // Action terminals embed the taskId mid-key (`action-setup--…`), + // so the tagged taskId is authoritative; the key match covers instances + // created without one. const key = instance.persistenceKey; - if (key === taskId || key.startsWith(`${taskId}-`)) { + if ( + instance.taskId === taskId || + key === taskId || + key.startsWith(`${taskId}-`) + ) { this.destroy(sessionId); } } From e51652c3eba2066498e4908ab063ab14b65ac228 Mon Sep 17 00:00:00 2001 From: Charles Vien Date: Wed, 1 Jul 2026 19:21:15 -0700 Subject: [PATCH 19/20] destroy terminals only after archive succeeds --- .../src/archive/archiveOrchestration.test.ts | 23 +++++++++++++++++-- .../core/src/archive/archiveOrchestration.ts | 10 +++----- .../ui/src/features/archive/useArchiveTask.ts | 18 --------------- 3 files changed, 24 insertions(+), 27 deletions(-) diff --git a/packages/core/src/archive/archiveOrchestration.test.ts b/packages/core/src/archive/archiveOrchestration.test.ts index f95370da1c..1058e19286 100644 --- a/packages/core/src/archive/archiveOrchestration.test.ts +++ b/packages/core/src/archive/archiveOrchestration.test.ts @@ -18,9 +18,7 @@ class Harness { unpin: vi.fn().mockResolvedValue(undefined), togglePin: vi.fn().mockResolvedValue(undefined), navigateAwayFromTaskIfActive: vi.fn(), - snapshotTerminalStates: vi.fn().mockReturnValue({}), clearTerminalStates: vi.fn(), - restoreTerminalStates: vi.fn(), snapshotCommandCenter: vi .fn() .mockReturnValue({ index: -1, wasActive: false }), @@ -103,6 +101,27 @@ describe("archiveTask", () => { expect(harness.list).toEqual([]); expect(harness.deps.togglePin).toHaveBeenCalledWith(TASK_ID); }); + + it("destroys terminals only after the archive succeeds", async () => { + let clearedWhenArchiveCalled = true; + harness.deps.archive = vi.fn().mockImplementation(async () => { + clearedWhenArchiveCalled = + vi.mocked(harness.deps.clearTerminalStates).mock.calls.length > 0; + }); + + await archiveTask(TASK_ID, harness.deps); + + expect(clearedWhenArchiveCalled).toBe(false); + expect(harness.deps.clearTerminalStates).toHaveBeenCalledWith(TASK_ID); + }); + + it("keeps terminals when archive fails", async () => { + harness.deps.archive = vi.fn().mockRejectedValue(new Error("boom")); + + await expect(archiveTask(TASK_ID, harness.deps)).rejects.toThrow("boom"); + + expect(harness.deps.clearTerminalStates).not.toHaveBeenCalled(); + }); }); describe("archiveTasks", () => { diff --git a/packages/core/src/archive/archiveOrchestration.ts b/packages/core/src/archive/archiveOrchestration.ts index a8d8406c41..b7e4c75ceb 100644 --- a/packages/core/src/archive/archiveOrchestration.ts +++ b/packages/core/src/archive/archiveOrchestration.ts @@ -27,9 +27,7 @@ export interface ArchiveOrchestrationDeps { unpin(taskId: string): Promise; togglePin(taskId: string): Promise; navigateAwayFromTaskIfActive(taskId: string): void; - snapshotTerminalStates(taskId: string): Record; clearTerminalStates(taskId: string): void; - restoreTerminalStates(states: Record): void; snapshotCommandCenter(taskId: string): { index: number; wasActive: boolean }; removeFromCommandCenter(taskId: string): void; restoreCommandCenter( @@ -69,11 +67,9 @@ export async function archiveTask( deps.navigateAwayFromTaskIfActive(taskId); } - const terminalStatesSnapshot = deps.snapshotTerminalStates(taskId); const commandCenterSnapshot = deps.snapshotCommandCenter(taskId); await deps.unpin(taskId); - deps.clearTerminalStates(taskId); deps.removeFromCommandCenter(taskId); await deps.cache.cancelPathFilter(); @@ -101,6 +97,9 @@ export async function archiveTask( try { await deps.disconnectFromTask(taskId); await deps.archive(taskId); + // Destroying terminals is irreversible, so it waits for the archive to + // commit; a failed archive keeps its live terminals. + deps.clearTerminalStates(taskId); // Non-optimistic flows keep the row visible during the request, then remove // it the moment the archive succeeds. if (!optimistic) { @@ -115,9 +114,6 @@ export async function archiveTask( if (wasPinned) { await deps.togglePin(taskId); } - if (Object.keys(terminalStatesSnapshot).length > 0) { - deps.restoreTerminalStates(terminalStatesSnapshot); - } if (commandCenterSnapshot.index !== -1) { deps.restoreCommandCenter(taskId, commandCenterSnapshot); } diff --git a/packages/ui/src/features/archive/useArchiveTask.ts b/packages/ui/src/features/archive/useArchiveTask.ts index 605c4aa576..9d1e60979a 100644 --- a/packages/ui/src/features/archive/useArchiveTask.ts +++ b/packages/ui/src/features/archive/useArchiveTask.ts @@ -20,10 +20,6 @@ import { useCommandCenterStore } from "@posthog/ui/features/command-center/comma import { useFocusStore } from "@posthog/ui/features/focus/focusStore"; import { pinnedTasksApi } from "@posthog/ui/features/sidebar/taskMetaApi"; import { destroyTaskTerminals } from "@posthog/ui/features/terminal/destroyTaskTerminals"; -import { - type TerminalState, - useTerminalStore, -} from "@posthog/ui/features/terminal/terminalStore"; import { toast } from "@posthog/ui/primitives/toast"; import { getAppViewSnapshot } from "@posthog/ui/router/useAppView"; import { openTaskInput } from "@posthog/ui/router/useOpenTask"; @@ -97,21 +93,7 @@ function makeOrchestrationDeps( ); } }, - snapshotTerminalStates: (taskId) => - Object.fromEntries( - Object.entries(useTerminalStore.getState().terminalStates).filter( - ([key]) => key === taskId || key.startsWith(`${taskId}-`), - ), - ), clearTerminalStates: (taskId) => destroyTaskTerminals(taskId), - restoreTerminalStates: (states) => { - useTerminalStore.setState((s) => ({ - terminalStates: { - ...s.terminalStates, - ...(states as Record), - }, - })); - }, snapshotCommandCenter: (taskId) => { const state = useCommandCenterStore.getState(); return { From f54fc25727bf15295e7123506633c70e753d9e1c Mon Sep 17 00:00:00 2001 From: Charles Vien Date: Wed, 1 Jul 2026 19:21:17 -0700 Subject: [PATCH 20/20] cover eviction protection and teardown seams --- .../core/src/sessions/sessionEviction.test.ts | 29 ++- .../sessions/sessionServiceEviction.test.ts | 185 ++++++++++++++++-- .../hooks/useSessionConnection.test.tsx | 94 +++++++++ .../tasks/useTaskCrudMutations.test.tsx | 46 ++++- .../terminal/destroyTaskTerminals.test.ts | 29 +++ 5 files changed, 365 insertions(+), 18 deletions(-) create mode 100644 packages/ui/src/features/sessions/hooks/useSessionConnection.test.tsx create mode 100644 packages/ui/src/features/terminal/destroyTaskTerminals.test.ts diff --git a/packages/core/src/sessions/sessionEviction.test.ts b/packages/core/src/sessions/sessionEviction.test.ts index 26fb8ee228..71f1ccb1c5 100644 --- a/packages/core/src/sessions/sessionEviction.test.ts +++ b/packages/core/src/sessions/sessionEviction.test.ts @@ -1,6 +1,11 @@ import type { AgentSession } from "@posthog/shared"; import { describe, expect, it } from "vitest"; -import { isSessionIdle, selectSessionsToEvict } from "./sessionEviction"; +import { getCellCount, type LayoutPreset } from "../command-center/grid"; +import { + isSessionIdle, + MAX_CONNECTED_SESSIONS, + selectSessionsToEvict, +} from "./sessionEviction"; function makeSession(overrides: Partial): AgentSession { return { @@ -19,6 +24,8 @@ describe("isSessionIdle", () => { ["connected idle local session", {}, true], ["connecting session", { status: "connecting" as const }, false], ["pending prompt", { isPromptPending: true }, false], + ["compacting session", { isCompacting: true }, false], + ["handoff in progress", { handoffInProgress: true }, false], [ "pending permission", { pendingPermissions: new Map([["p1", {} as never]]) }, @@ -113,3 +120,23 @@ describe("selectSessionsToEvict", () => { expect(evicted.map((s) => s.taskId)).toEqual(expected); }); }); + +describe("MAX_CONNECTED_SESSIONS", () => { + it("stays above the largest Command Center grid so full layouts never evict", () => { + // Record forces this list to grow with the union, so a + // new larger preset breaks this test instead of silently churning cells. + const allPresets: Record = { + "1x1": true, + "2x1": true, + "1x2": true, + "2x2": true, + "3x2": true, + "3x3": true, + }; + const largestGrid = Math.max( + ...(Object.keys(allPresets) as LayoutPreset[]).map(getCellCount), + ); + + expect(MAX_CONNECTED_SESSIONS).toBeGreaterThan(largestGrid); + }); +}); diff --git a/packages/core/src/sessions/sessionServiceEviction.test.ts b/packages/core/src/sessions/sessionServiceEviction.test.ts index ce329977c2..081373edfe 100644 --- a/packages/core/src/sessions/sessionServiceEviction.test.ts +++ b/packages/core/src/sessions/sessionServiceEviction.test.ts @@ -1,8 +1,12 @@ import type { AgentSession } from "@posthog/shared"; import type { Task } from "@posthog/shared/domain-types"; -import { describe, expect, it, vi } from "vitest"; +import { afterEach, describe, expect, it, vi } from "vitest"; import { MAX_CONNECTED_SESSIONS } from "./sessionEviction"; -import { SessionService, type SessionServiceDeps } from "./sessionService"; +import { + type ReconcileTaskConnectionParams, + SessionService, + type SessionServiceDeps, +} from "./sessionService"; function makeSession( taskId: string, @@ -37,15 +41,18 @@ function createHarness(seedSessions: AgentSession[]) { delete sessions[taskRunId]; }); const cancelMutate = vi.fn().mockResolvedValue(undefined); + const removePersistedConfigOptions = vi.fn(); + const removeAdapter = vi.fn(); const store = { getSessions: () => sessions, getSessionByTaskId: (taskId: string) => Object.values(sessions).find((s) => s.taskId === taskId), removeSession, + updateSession: vi.fn(), }; - const noopLog = { + const log = { info: vi.fn(), warn: vi.fn(), error: vi.fn(), @@ -54,14 +61,14 @@ function createHarness(seedSessions: AgentSession[]) { const deps = { store, - log: noopLog, + log, getPersistedConfigOptions: () => undefined, setPersistedConfigOptions: vi.fn(), - removePersistedConfigOptions: vi.fn(), + removePersistedConfigOptions, adapterStore: { getAdapter: () => undefined, setAdapter: vi.fn(), - removeAdapter: vi.fn(), + removeAdapter, }, trpc: { agent: { @@ -74,7 +81,15 @@ function createHarness(seedSessions: AgentSession[]) { } as unknown as SessionServiceDeps; const service = new SessionService(deps); - return { service, sessions, removeSession, cancelMutate }; + return { + service, + sessions, + removeSession, + cancelMutate, + removePersistedConfigOptions, + removeAdapter, + log, + }; } function connectParamsFor(taskId: string) { @@ -84,12 +99,19 @@ function connectParamsFor(taskId: string) { }; } +function seedIdleSessions(prefix = "idle", count = MAX_CONNECTED_SESSIONS) { + return Array.from({ length: count }, (_, i) => + makeSession(`${prefix}-${i}`, i + 1), + ); +} + describe("SessionService idle session eviction", () => { + afterEach(() => { + vi.useRealTimers(); + }); + it("evicts the least recently used idle sessions beyond the budget", async () => { - const idleCount = MAX_CONNECTED_SESSIONS; - const seeds = Array.from({ length: idleCount }, (_, i) => - makeSession(`idle-${i}`, i + 1), - ); + const seeds = seedIdleSessions(); seeds.push(makeSession("active", 1000)); const { service, removeSession } = createHarness(seeds); @@ -103,8 +125,7 @@ describe("SessionService idle session eviction", () => { }); it("never evicts mounted or busy sessions", async () => { - const idleCount = MAX_CONNECTED_SESSIONS; - const seeds = Array.from({ length: idleCount }, (_, i) => + const seeds = Array.from({ length: MAX_CONNECTED_SESSIONS }, (_, i) => makeSession(`idle-${i}`, i + 1, { isPromptPending: i === 1, }), @@ -127,9 +148,7 @@ describe("SessionService idle session eviction", () => { }); it("evicts nothing at or under the budget", async () => { - const seeds = Array.from({ length: MAX_CONNECTED_SESSIONS - 2 }, (_, i) => - makeSession(`idle-${i}`, i + 1), - ); + const seeds = seedIdleSessions("idle", MAX_CONNECTED_SESSIONS - 2); seeds.push(makeSession("active", 1000)); const { service, removeSession } = createHarness(seeds); @@ -138,4 +157,138 @@ describe("SessionService idle session eviction", () => { expect(removeSession).not.toHaveBeenCalled(); }); + + it("preserves the adapter and persisted config of evicted sessions", async () => { + const seeds = seedIdleSessions(); + seeds.push(makeSession("active", 1000)); + const { + service, + removeSession, + removePersistedConfigOptions, + removeAdapter, + } = createHarness(seeds); + + await service.connectToTask(connectParamsFor("active")); + + await vi.waitFor(() => { + expect(removeSession).toHaveBeenCalledTimes(2); + }); + expect(removePersistedConfigOptions).not.toHaveBeenCalled(); + expect(removeAdapter).not.toHaveBeenCalled(); + }); + + it("drops the adapter and persisted config on a real disconnect", async () => { + const { service, removePersistedConfigOptions, removeAdapter } = + createHarness([makeSession("idle-0", 1)]); + + await service.disconnectFromTask("idle-0"); + + expect(removePersistedConfigOptions).toHaveBeenCalledWith("run-idle-0"); + expect(removeAdapter).toHaveBeenCalledWith("run-idle-0"); + }); + + it("bounds the budget when reconciling a cloud task", async () => { + const seeds = seedIdleSessions(); + seeds.push(makeSession("cloud-active", 1000, { isCloud: true })); + const { service, removeSession } = createHarness(seeds); + + service.reconcileTaskConnection({ + task: { + id: "cloud-active", + title: "cloud-active", + description: "cloud-active", + } as Task, + session: undefined, + repoPath: null, + isCloud: true, + isOnline: true, + cloudAuth: { status: "loading" }, + } as ReconcileTaskConnectionParams); + + await vi.waitFor(() => { + expect(removeSession).toHaveBeenCalledTimes(2); + }); + expect(removeSession).toHaveBeenCalledWith("run-idle-0"); + expect(removeSession).not.toHaveBeenCalledWith("run-cloud-active"); + }); + + it("evicts cloud sessions once their runs are terminal", async () => { + const seeds = Array.from({ length: MAX_CONNECTED_SESSIONS }, (_, i) => + makeSession(`cloud-${i}`, i + 1, { + isCloud: true, + cloudStatus: "completed", + }), + ); + seeds.push(makeSession("active", 1000)); + const { service, removeSession } = createHarness(seeds); + + await service.connectToTask(connectParamsFor("active")); + + await vi.waitFor(() => { + expect(removeSession).toHaveBeenCalledTimes(2); + }); + expect(removeSession).toHaveBeenCalledWith("run-cloud-0"); + expect(removeSession).toHaveBeenCalledWith("run-cloud-1"); + }); + + it("keeps a task protected while any mount remains", async () => { + const seeds = seedIdleSessions(); + seeds.push(makeSession("active", 1000)); + const { service, removeSession } = createHarness(seeds); + + const unregisterFirst = service.registerMountedTask("idle-0"); + service.registerMountedTask("idle-0"); + unregisterFirst(); + + await service.connectToTask(connectParamsFor("active")); + + await vi.waitFor(() => { + expect(removeSession).toHaveBeenCalledTimes(2); + }); + expect(removeSession).not.toHaveBeenCalledWith("run-idle-0"); + }); + + it("frees a task for eviction after its last unmount", async () => { + vi.useFakeTimers({ toFake: ["Date"] }); + const seeds = seedIdleSessions(); + seeds.push(makeSession("active", 1000)); + const { service, removeSession } = createHarness(seeds); + + vi.setSystemTime(1_000); + const unregisterFirst = service.registerMountedTask("idle-0"); + const unregisterSecond = service.registerMountedTask("idle-0"); + unregisterFirst(); + unregisterSecond(); + + vi.setSystemTime(2_000); + for (let i = 1; i < MAX_CONNECTED_SESSIONS; i++) { + service.registerMountedTask(`idle-${i}`)(); + } + + vi.setSystemTime(3_000); + await service.connectToTask(connectParamsFor("active")); + + await vi.waitFor(() => { + expect(removeSession).toHaveBeenCalledWith("run-idle-0"); + }); + }); + + it("keeps evicting after one teardown fails", async () => { + const seeds = seedIdleSessions(); + seeds.push(makeSession("active", 1000)); + const { service, removeSession, log } = createHarness(seeds); + removeSession.mockImplementationOnce(() => { + throw new Error("dispose failed"); + }); + + await service.connectToTask(connectParamsFor("active")); + + await vi.waitFor(() => { + expect(removeSession).toHaveBeenCalledWith("run-idle-1"); + }); + expect(log.error).toHaveBeenCalledWith( + "Failed to evict idle session", + expect.objectContaining({ taskId: "idle-0" }), + ); + }); }); diff --git a/packages/ui/src/features/sessions/hooks/useSessionConnection.test.tsx b/packages/ui/src/features/sessions/hooks/useSessionConnection.test.tsx new file mode 100644 index 0000000000..4c3d163022 --- /dev/null +++ b/packages/ui/src/features/sessions/hooks/useSessionConnection.test.tsx @@ -0,0 +1,94 @@ +import type { Task } from "@posthog/shared/domain-types"; +import { renderHook } from "@testing-library/react"; +import { beforeEach, describe, expect, it, vi } from "vitest"; + +const mocks = vi.hoisted(() => { + const unregisterMountedTask = vi.fn(); + const sessionService = { + registerMountedTask: vi.fn(() => unregisterMountedTask), + startActivityHeartbeat: vi.fn(() => () => {}), + reconcileTaskConnection: vi.fn(() => () => {}), + }; + return { sessionService, unregisterMountedTask }; +}); + +vi.mock("@posthog/di/react", () => ({ + useService: () => mocks.sessionService, +})); + +vi.mock("@tanstack/react-query", () => ({ + useQueryClient: () => ({ invalidateQueries: vi.fn() }), +})); + +vi.mock("@posthog/ui/hooks/useConnectivity", () => ({ + useConnectivity: () => ({ isOnline: true }), +})); + +vi.mock("@posthog/ui/features/auth/store", () => ({ + useAuthStateValue: (selector: (state: Record) => unknown) => + selector({ + status: "unauthenticated", + bootstrapComplete: false, + currentProjectId: null, + cloudRegion: null, + }), +})); + +vi.mock("./useChatTitleGenerator", () => ({ + useChatTitleGenerator: vi.fn(), +})); + +import { useSessionConnection } from "./useSessionConnection"; + +function makeTask(id: string): Task { + return { id, title: id, description: id } as Task; +} + +function connectionProps(taskId: string) { + return { + taskId, + task: makeTask(taskId), + session: undefined, + repoPath: null, + isCloud: false, + }; +} + +describe("useSessionConnection mounted-task registration", () => { + beforeEach(() => { + mocks.sessionService.registerMountedTask.mockClear(); + mocks.unregisterMountedTask.mockClear(); + }); + + it("registers the task on mount and unregisters on unmount", () => { + const { unmount } = renderHook(() => + useSessionConnection(connectionProps("task-1")), + ); + + expect(mocks.sessionService.registerMountedTask).toHaveBeenCalledWith( + "task-1", + ); + expect(mocks.unregisterMountedTask).not.toHaveBeenCalled(); + + unmount(); + expect(mocks.unregisterMountedTask).toHaveBeenCalledTimes(1); + }); + + it("re-registers when the task changes", () => { + const { rerender, unmount } = renderHook( + ({ taskId }: { taskId: string }) => + useSessionConnection(connectionProps(taskId)), + { initialProps: { taskId: "task-1" } }, + ); + + rerender({ taskId: "task-2" }); + + expect(mocks.unregisterMountedTask).toHaveBeenCalledTimes(1); + expect(mocks.sessionService.registerMountedTask).toHaveBeenLastCalledWith( + "task-2", + ); + + unmount(); + expect(mocks.unregisterMountedTask).toHaveBeenCalledTimes(2); + }); +}); diff --git a/packages/ui/src/features/tasks/useTaskCrudMutations.test.tsx b/packages/ui/src/features/tasks/useTaskCrudMutations.test.tsx index 7270003412..76bde2513e 100644 --- a/packages/ui/src/features/tasks/useTaskCrudMutations.test.tsx +++ b/packages/ui/src/features/tasks/useTaskCrudMutations.test.tsx @@ -19,12 +19,22 @@ const confirmAndDelete = vi.hoisted(() => const deletionService = vi.hoisted(() => ({ deleteTask: vi.fn().mockResolvedValue(undefined), confirmAndDelete, + disconnectFromTask: vi.fn().mockResolvedValue(undefined), })); const destroyTaskTerminals = vi.hoisted(() => vi.fn()); +const capturedMutationFns = vi.hoisted( + () => [] as Array<(client: unknown, vars: string) => Promise>, +); + vi.mock("@posthog/ui/hooks/useAuthenticatedMutation", () => ({ - useAuthenticatedMutation: () => ({ mutateAsync, isPending: false }), + useAuthenticatedMutation: ( + fn: (client: unknown, vars: string) => Promise, + ) => { + capturedMutationFns.push(fn); + return { mutateAsync, isPending: false }; + }, })); vi.mock("@posthog/di/react", () => ({ useService: () => deletionService, @@ -98,6 +108,40 @@ describe("useDeleteTask.deleteWithConfirm", () => { }); }); +describe("useDeleteTask mutation", () => { + beforeEach(() => { + vi.clearAllMocks(); + capturedMutationFns.length = 0; + }); + + it("releases resources only after the server delete succeeds", async () => { + deletionService.deleteTask.mockResolvedValueOnce("deleted"); + renderHook(() => useDeleteTask(), { wrapper }); + + const mutationFn = capturedMutationFns.at(-1); + const result = await mutationFn?.("client", "t1"); + + expect(result).toBe("deleted"); + expect(deletionService.deleteTask).toHaveBeenCalledWith("client", "t1"); + expect(deletionService.disconnectFromTask).toHaveBeenCalledWith("t1"); + expect(destroyTaskTerminals).toHaveBeenCalledWith("t1"); + expect(deletionService.deleteTask.mock.invocationCallOrder[0]).toBeLessThan( + deletionService.disconnectFromTask.mock.invocationCallOrder[0], + ); + }); + + it("does not release resources when the server delete fails", async () => { + deletionService.deleteTask.mockRejectedValueOnce(new Error("nope")); + renderHook(() => useDeleteTask(), { wrapper }); + + const mutationFn = capturedMutationFns.at(-1); + await expect(mutationFn?.("client", "t1")).rejects.toThrow("nope"); + + expect(deletionService.disconnectFromTask).not.toHaveBeenCalled(); + expect(destroyTaskTerminals).not.toHaveBeenCalled(); + }); +}); + describe("releaseDeletedTaskResources", () => { beforeEach(() => { vi.clearAllMocks(); diff --git a/packages/ui/src/features/terminal/destroyTaskTerminals.test.ts b/packages/ui/src/features/terminal/destroyTaskTerminals.test.ts new file mode 100644 index 0000000000..d48ec515ee --- /dev/null +++ b/packages/ui/src/features/terminal/destroyTaskTerminals.test.ts @@ -0,0 +1,29 @@ +import { describe, expect, it, vi } from "vitest"; + +const mocks = vi.hoisted(() => ({ + destroyForTask: vi.fn(), + clearTerminalStatesForTask: vi.fn(), +})); + +vi.mock("./TerminalManager", () => ({ + terminalManager: { destroyForTask: mocks.destroyForTask }, +})); + +vi.mock("./terminalStore", () => ({ + useTerminalStore: { + getState: () => ({ + clearTerminalStatesForTask: mocks.clearTerminalStatesForTask, + }), + }, +})); + +import { destroyTaskTerminals } from "./destroyTaskTerminals"; + +describe("destroyTaskTerminals", () => { + it("destroys live instances and clears persisted state for the task", () => { + destroyTaskTerminals("task-1"); + + expect(mocks.destroyForTask).toHaveBeenCalledWith("task-1"); + expect(mocks.clearTerminalStatesForTask).toHaveBeenCalledWith("task-1"); + }); +});