From 2fc3b7b15efca55cefcb3e07e254d4750cf46c2f Mon Sep 17 00:00:00 2001 From: Merge Sim Date: Thu, 3 Sep 2026 00:35:08 -0700 Subject: [PATCH] fix(wsl): coalesce relay identity reads across polling --- src/main/agent-hooks/server.ts | 4 + .../server/server-ingest-remote.ts | 14 ++ .../wsl-hook-relay-guest-install.ts | 13 ++ .../agent-hooks/wsl-hook-relay-identity.ts | 2 +- .../wsl-hook-relay-manager.test.ts | 7 + .../agent-hooks/wsl-hook-relay-manager.ts | 8 +- .../daemon/daemon-pty-session-inventory.ts | 1 + src/main/daemon/daemon-pty-session-spawn.ts | 2 + ...cal-pty-provider-session-inventory.test.ts | 39 +++++ .../providers/local-pty-session-activation.ts | 10 ++ .../providers/local-pty-session-operations.ts | 3 + .../wsl-relay-identity-reader.test.ts | 163 ++++++++++++++++++ .../providers/wsl-relay-identity-reader.ts | 60 +++++-- ...refresh-floating-workspace-pty-liveness.ts | 12 +- ...sh-pty-worktree-records-from-controller.ts | 36 +++- ...ktree-records-with-controller-inventory.ts | 24 +-- src/main/runtime/orca-runtime-runtime-id.ts | 26 ++- .../mobile-summaries-part-04.spec.ts | 151 ++++++++++++++++ 18 files changed, 541 insertions(+), 34 deletions(-) diff --git a/src/main/agent-hooks/server.ts b/src/main/agent-hooks/server.ts index f03484f14f9..75dfd94bdb3 100644 --- a/src/main/agent-hooks/server.ts +++ b/src/main/agent-hooks/server.ts @@ -22,6 +22,10 @@ export { RETIRED_PANE_FENCES_MAX } from './server/server-constants' export { isValidPaneKey } +export { + getAgentHookRemotePostCount, + resetAgentHookRemotePostCount +} from './server/server-ingest-remote' /** Public composition seam for the loopback hook listener and relay status adapter. */ export class AgentHookServer extends AgentHookServerLifecycle {} diff --git a/src/main/agent-hooks/server/server-ingest-remote.ts b/src/main/agent-hooks/server/server-ingest-remote.ts index 18b0a837c32..c6ea1a5152f 100644 --- a/src/main/agent-hooks/server/server-ingest-remote.ts +++ b/src/main/agent-hooks/server/server-ingest-remote.ts @@ -19,6 +19,19 @@ import type { AgentHookEventPayload } from '../../../shared/agent-hook-listener/ import { isValidPiProviderSessionOnly } from './server-status-identity' import { AgentHookServerIngestTerminal } from './server-ingest-terminal' +// Diagnostics seam for the WSL hooks-off contract. This counts relay posts +// reaching the host boundary; local loopback hook traffic is intentionally not +// included. +let agentHookRemotePostCount = 0 + +export function getAgentHookRemotePostCount(): number { + return agentHookRemotePostCount +} + +export function resetAgentHookRemotePostCount(): void { + agentHookRemotePostCount = 0 +} + export abstract class AgentHookServerIngestRemote extends AgentHookServerIngestTerminal { /** Ingest a payload from the relay JSON-RPC channel (not the local HTTP server); connectionId is stamped here. Main is still the SSH trust boundary, so re-run the canonical normalizer before caching. */ ingestRemote( @@ -49,6 +62,7 @@ export abstract class AgentHookServerIngestRemote extends AgentHookServerIngestT }, connectionId: string | null ): void { + agentHookRemotePostCount++ // Why: wire crosses a trust boundary — re-check/trim so an empty connectionId can't poison caches. if (connectionId !== null && typeof connectionId !== 'string') { return diff --git a/src/main/agent-hooks/wsl-hook-relay-guest-install.ts b/src/main/agent-hooks/wsl-hook-relay-guest-install.ts index f356e7c59f9..a916859f7c6 100644 --- a/src/main/agent-hooks/wsl-hook-relay-guest-install.ts +++ b/src/main/agent-hooks/wsl-hook-relay-guest-install.ts @@ -31,6 +31,18 @@ type GuestInstallState = { lastInstallAt?: number } +// Diagnostics seam for the hooks-off contract: relay residency must remain +// live while no guest hook installation work is attempted. +let wslRelayGuestHookInstallCount = 0 + +export function getWslRelayGuestHookInstallCount(): number { + return wslRelayGuestHookInstallCount +} + +export function resetWslRelayGuestHookInstallCount(): void { + wslRelayGuestHookInstallCount = 0 +} + export async function runWslRelayGuestInstall( deps: GuestInstallDeps, state: GuestInstallState, @@ -43,6 +55,7 @@ export async function runWslRelayGuestInstall( ) { return } + wslRelayGuestHookInstallCount++ state.lastInstallAt = Date.now() await installWslGuestHooks({ mux, diff --git a/src/main/agent-hooks/wsl-hook-relay-identity.ts b/src/main/agent-hooks/wsl-hook-relay-identity.ts index 64c9475459d..f7b8f8ebffb 100644 --- a/src/main/agent-hooks/wsl-hook-relay-identity.ts +++ b/src/main/agent-hooks/wsl-hook-relay-identity.ts @@ -19,7 +19,7 @@ export async function readWslRelayProcessIdentity(options: { ensure: (distro: string) => Promise | void getState: (distro: string) => { phase: string; mux?: SshChannelMultiplexer } | undefined disposed: boolean - requestOptions?: { signal?: AbortSignal; timeoutMs?: number } + requestOptions?: { signal?: AbortSignal; timeoutMs?: number; stableUntilReset?: boolean } }): Promise { const unavailable = (reason: string): WslRelayIdentityResult[] => options.anchors.map(() => ({ status: 'unverifiable' as const, reason, capturedAgeMs: 0 })) diff --git a/src/main/agent-hooks/wsl-hook-relay-manager.test.ts b/src/main/agent-hooks/wsl-hook-relay-manager.test.ts index 09bfebfe2a9..147f37891c9 100644 --- a/src/main/agent-hooks/wsl-hook-relay-manager.test.ts +++ b/src/main/agent-hooks/wsl-hook-relay-manager.test.ts @@ -18,8 +18,11 @@ import { import { getWslRelayIdentityRpcCount, resetWslRelayIdentityRpcCount, + getWslRelayGuestHookInstallCount, + resetWslRelayGuestHookInstallCount, WslHookRelayManager } from './wsl-hook-relay-manager' +import { getAgentHookRemotePostCount, resetAgentHookRemotePostCount } from './server' import { FAILURE_COOLDOWN_BASE_MS, type WslHookRelayManagerDeps } from './wsl-hook-relay-deps' import { AGENT_HOOK_INSTALL_PLUGINS_METHOD, @@ -450,6 +453,8 @@ describe('WslHookRelayManager', () => { it('serves identity reads through connectWslRelayState with hooks disabled', async () => { resetWslRelayIdentityRpcCount() + resetWslRelayGuestHookInstallCount() + resetAgentHookRemotePostCount() const waitForSentinel = vi.fn(async () => { const transport = guestTransport() const guest = harnesses.at(-1)!.guestDispatcher @@ -482,6 +487,8 @@ describe('WslHookRelayManager', () => { { status: 'live', processName: 'claude' } ]) expect(getWslRelayIdentityRpcCount()).toBe(1) + expect(getWslRelayGuestHookInstallCount()).toBe(0) + expect(getAgentHookRemotePostCount()).toBe(0) expect(deps.installHooks).not.toHaveBeenCalled() manager.disposeAll() }) diff --git a/src/main/agent-hooks/wsl-hook-relay-manager.ts b/src/main/agent-hooks/wsl-hook-relay-manager.ts index 60617683c67..d370f9cf7db 100644 --- a/src/main/agent-hooks/wsl-hook-relay-manager.ts +++ b/src/main/agent-hooks/wsl-hook-relay-manager.ts @@ -24,6 +24,10 @@ import { wslRelayIdentityRpcCount } from './wsl-hook-relay-identity' import { wslHookRelayStateKey } from './wsl-hook-relay-state-key' +import { + getWslRelayGuestHookInstallCount, + resetWslRelayGuestHookInstallCount +} from './wsl-hook-relay-guest-install' import { sanitizeWslHookInstanceKey, WSL_RELAY_HOOKS_SET_ENABLED_METHOD @@ -44,6 +48,8 @@ export function resetWslRelayIdentityRpcCount(): void { resetIdentityCounter() } +export { getWslRelayGuestHookInstallCount, resetWslRelayGuestHookInstallCount } + export class WslHookRelayManager { private deps: WslHookRelayManagerDeps private recovery: WslRelayRecovery @@ -155,7 +161,7 @@ export class WslHookRelayManager { readProcessIdentity = ( distro: string, anchors: Parameters[0]['anchors'], - options?: { signal?: AbortSignal; timeoutMs?: number } + options?: { signal?: AbortSignal; timeoutMs?: number; stableUntilReset?: boolean } ) => readWslRelayProcessIdentity({ distro, diff --git a/src/main/daemon/daemon-pty-session-inventory.ts b/src/main/daemon/daemon-pty-session-inventory.ts index f73efdd6b52..4ad8a66b945 100644 --- a/src/main/daemon/daemon-pty-session-inventory.ts +++ b/src/main/daemon/daemon-pty-session-inventory.ts @@ -71,6 +71,7 @@ export abstract class DaemonPtySessionInventory extends DaemonPtyProcessInspecti sessions.map((session) => session.wslShellAnchor!), { signal: opts?.signal, + stableUntilReset: true, timeoutMs: opts?.deadlineMs === undefined ? undefined diff --git a/src/main/daemon/daemon-pty-session-spawn.ts b/src/main/daemon/daemon-pty-session-spawn.ts index 235b286b0bc..2b9137a6435 100644 --- a/src/main/daemon/daemon-pty-session-spawn.ts +++ b/src/main/daemon/daemon-pty-session-spawn.ts @@ -20,12 +20,14 @@ import { resolveWslSessionContext } from './wsl-session-context' import { resolveSafePtyDefaultCwd } from '../providers/pty-default-cwd' import { resolveUnixShellPath } from '../providers/local-pty-utils' import type { PtySpawnOptions, PtySpawnResult } from '../providers/types' +import { wslRelayIdentityReader } from '../providers/wsl-relay-identity-reader' import { injectHistoryEnv, injectWslFishHistoryEnv, logHistoryInjection } from '../terminal-history' import { addWslEnvKeys } from '../wsl-env' export abstract class DaemonPtySessionSpawn extends DaemonPtySpawnResult { async spawn(opts: PtySpawnOptions): Promise { const spawnOpts = this.withHistoryIsolation(opts) + wslRelayIdentityReader.reset() const sessionId = spawnOpts.sessionId ?? mintPtySessionId(spawnOpts.worktreeId) const operation: PendingDaemonSpawnOperation = { exitsBySessionId: new Map(), diff --git a/src/main/providers/local-pty-provider-session-inventory.test.ts b/src/main/providers/local-pty-provider-session-inventory.test.ts index 3c6f069ba47..773d01d54b1 100644 --- a/src/main/providers/local-pty-provider-session-inventory.test.ts +++ b/src/main/providers/local-pty-provider-session-inventory.test.ts @@ -122,6 +122,7 @@ import { ptyWslShellAnchors } from './local-pty-provider-state' import { wslRelayIdentityReader } from './wsl-relay-identity-reader' +import { wslHookRelayManager } from '../agent-hooks/wsl-hook-relay-manager' import { applyLocalPtyProviderMockDefaults, createLocalPtyMockProcess, @@ -250,6 +251,44 @@ describe('LocalPtyProvider', () => { readBatch.mockRestore() } }) + + it('coalesces repeated WSL list-process bursts until a PTY event resets identity', async () => { + Object.defineProperty(process, 'platform', { configurable: true, value: 'win32' }) + const spawned = await provider.spawn({ cols: 80, rows: 24, cwd: '/tmp/wsl-owned-cwd' }) + const anchor = { + distro: 'Ubuntu', + bootId: '11111111-1111-1111-1111-111111111111', + shellPid: 100, + shellStartTime: 1, + tty: '/dev/pts/1' + } + ptyWslDistroById.set(spawned.id, 'Ubuntu') + ptyWslShellAnchors.set(spawned.id, anchor) + wslRelayIdentityReader.reset() + const identityRequest = vi + .spyOn(wslHookRelayManager, 'readProcessIdentity') + .mockResolvedValue([ + { + status: 'unverifiable' as const, + reason: 'relay_unavailable', + capturedAgeMs: 0 + } + ]) + try { + await Promise.all(Array.from({ length: 5 }, () => provider.listProcesses())) + expect(identityRequest).toHaveBeenCalledOnce() + + await Promise.all(Array.from({ length: 5 }, () => provider.listProcesses())) + expect(identityRequest).toHaveBeenCalledOnce() + + wslRelayIdentityReader.reset() + await provider.listProcesses() + expect(identityRequest).toHaveBeenCalledTimes(2) + } finally { + identityRequest.mockRestore() + clearPtyState(spawned.id) + } + }) }) describe('getDefaultShell', () => { diff --git a/src/main/providers/local-pty-session-activation.ts b/src/main/providers/local-pty-session-activation.ts index ce78364f5e2..650e2bd84e0 100644 --- a/src/main/providers/local-pty-session-activation.ts +++ b/src/main/providers/local-pty-session-activation.ts @@ -41,6 +41,7 @@ import { type ShellStartupIdentityScanState } from '../shell-startup-identity-scanner' import type { PtySpawnOptions, PtySpawnResult } from './types' +import { wslRelayIdentityReader } from './wsl-relay-identity-reader' export function activateLocalPtySession(args: { id: string @@ -55,6 +56,8 @@ export function activateLocalPtySession(args: { }): PtySpawnResult { const { id, incarnationId, spawn, getOptions, plan, env, proc, spawnedWslDistro } = args createPtyPhysicalExit(id) + // A newly created PTY can introduce a fresh WSL shell identity. + wslRelayIdentityReader.reset() ptyReportsChildExitStatus.set(id, args.reportsChildExitStatus) ptyProcesses.set(id, proc) ptyInitialCwd.set(id, plan.cwd) @@ -78,13 +81,20 @@ export function activateLocalPtySession(args: { ) ptyLoadGeneration.set(id, getLoadGeneration()) ptyIncarnations.set(id, incarnationId) + // Foreground identity is event-invalidated: once output arrives, the next + // inventory may observe an agent transition instead of reusing the old read. + const resetWslIdentityCache = (): void => { + wslRelayIdentityReader.reset() + } getOptions().onSpawned?.(id, incarnationId) const emitIngressData = (emission: PtyIngressEmission): void => { const sequenceChars = emission.rawEndSeq - emission.rawStartSeq if (emission.transformed || sequenceChars !== emission.data.length) { + resetWslIdentityCache() getOptions().onData?.(id, emission.data, Date.now(), sequenceChars, true) } else { + resetWslIdentityCache() getOptions().onData?.(id, emission.data, Date.now()) } for (const cb of dataListeners) { diff --git a/src/main/providers/local-pty-session-operations.ts b/src/main/providers/local-pty-session-operations.ts index 397004ddacc..275fd2dd626 100644 --- a/src/main/providers/local-pty-session-operations.ts +++ b/src/main/providers/local-pty-session-operations.ts @@ -148,6 +148,9 @@ export async function listLocalPtyProcesses(opts?: { .filter((anchor): anchor is NonNullable => anchor !== undefined) const results = await wslRelayIdentityReader.readBatch(distro, anchors, { signal: opts?.signal, + // A process inventory is reused until a local PTY event resets the + // reader; polling itself is not evidence that guest state changed. + stableUntilReset: true, timeoutMs: opts?.deadlineMs === undefined ? undefined : Math.max(1, opts.deadlineMs - Date.now()) }) diff --git a/src/main/providers/wsl-relay-identity-reader.test.ts b/src/main/providers/wsl-relay-identity-reader.test.ts index 1220bd0620f..ae5920c1965 100644 --- a/src/main/providers/wsl-relay-identity-reader.test.ts +++ b/src/main/providers/wsl-relay-identity-reader.test.ts @@ -1,6 +1,7 @@ import { afterEach, describe, expect, it, vi } from 'vitest' import { createWslRelayIdentityReader } from './wsl-relay-identity-reader' +import type { WslRelayIdentityResult } from '../../shared/wsl-hook-relay-contract' describe('wsl relay identity reader cost invariant', () => { afterEach(() => { @@ -31,4 +32,166 @@ describe('wsl relay identity reader cost invariant', () => { await expect(reader.read('Ubuntu', anchor)).resolves.toMatchObject({ capturedAgeMs: 260 }) expect(readProcessIdentity).toHaveBeenCalledOnce() }) + + it('keeps inventory reads stable until a PTY event resets the cache', async () => { + vi.useFakeTimers() + vi.setSystemTime(0) + const readProcessIdentity = vi.fn().mockResolvedValue([ + { + status: 'live' as const, + processName: 'claude', + anchor: { + distro: 'Ubuntu', + bootId: '11111111-1111-1111-1111-111111111111', + shellPid: 1, + shellStartTime: 1, + tty: '/dev/pts/1' + }, + capturedAgeMs: 0 + } + ]) + const reader = createWslRelayIdentityReader({ readProcessIdentity }) + const anchor = { + distro: 'Ubuntu', + bootId: '11111111-1111-1111-1111-111111111111', + shellPid: 1, + shellStartTime: 1, + tty: '/dev/pts/1' + } + + await reader.readBatch('Ubuntu', [anchor], { stableUntilReset: true }) + vi.setSystemTime(60_000) + await reader.readBatch('Ubuntu', [anchor], { stableUntilReset: true }) + expect(readProcessIdentity).toHaveBeenCalledOnce() + + reader.reset() + await reader.readBatch('Ubuntu', [anchor], { stableUntilReset: true }) + expect(readProcessIdentity).toHaveBeenCalledTimes(2) + }) + + it('coalesces list-session bursts and reuses the capture until reset', async () => { + vi.useFakeTimers() + vi.setSystemTime(0) + const readProcessIdentity = vi.fn().mockResolvedValue([ + { + status: 'live' as const, + processName: 'claude', + anchor: { + distro: 'Ubuntu', + bootId: '11111111-1111-1111-1111-111111111111', + shellPid: 1, + shellStartTime: 1, + tty: '/dev/pts/1' + }, + capturedAgeMs: 0 + } + ]) + const reader = createWslRelayIdentityReader({ readProcessIdentity }) + const anchor = { + distro: 'Ubuntu', + bootId: '11111111-1111-1111-1111-111111111111', + shellPid: 1, + shellStartTime: 1, + tty: '/dev/pts/1' + } + + await Promise.all( + Array.from({ length: 5 }, () => + reader.readBatch('Ubuntu', [anchor], { stableUntilReset: true }) + ) + ) + expect(readProcessIdentity).toHaveBeenCalledOnce() + + vi.setSystemTime(60_000) + await Promise.all( + Array.from({ length: 5 }, () => + reader.readBatch('Ubuntu', [anchor], { stableUntilReset: true }) + ) + ) + expect(readProcessIdentity).toHaveBeenCalledOnce() + + reader.reset() + await reader.readBatch('Ubuntu', [anchor], { stableUntilReset: true }) + expect(readProcessIdentity).toHaveBeenCalledTimes(2) + }) + + it('does not republish an in-flight capture after an event reset', async () => { + let resolveCapture!: (result: WslRelayIdentityResult[]) => void + const replacementResult: WslRelayIdentityResult[] = [ + { + status: 'live', + processName: 'claude', + anchor: { + distro: 'Ubuntu', + bootId: '11111111-1111-1111-1111-111111111111', + shellPid: 1, + shellStartTime: 1, + tty: '/dev/pts/1' + }, + capturedAgeMs: 0 + } + ] + const readProcessIdentity = vi.fn( + () => + new Promise((resolve) => { + resolveCapture = resolve + }) + ) + const reader = createWslRelayIdentityReader({ readProcessIdentity }) + const anchor = { + distro: 'Ubuntu', + bootId: '11111111-1111-1111-1111-111111111111', + shellPid: 1, + shellStartTime: 1, + tty: '/dev/pts/1' + } + const first = reader.readBatch('Ubuntu', [anchor], { stableUntilReset: true }) + expect(readProcessIdentity).toHaveBeenCalledOnce() + reader.reset() + resolveCapture(replacementResult) + await first + + const second = reader.readBatch('Ubuntu', [anchor], { stableUntilReset: true }) + // The second request is intentionally independent of the stale first one. + resolveCapture(replacementResult) + await second + expect(readProcessIdentity).toHaveBeenCalledTimes(2) + }) + + it('keeps a stable snapshot marked stable after a foreground read', async () => { + vi.useFakeTimers() + vi.setSystemTime(0) + const readProcessIdentity = vi.fn().mockResolvedValue([ + { + status: 'live' as const, + processName: 'claude', + anchor: { + distro: 'Ubuntu', + bootId: '11111111-1111-1111-1111-111111111111', + shellPid: 1, + shellStartTime: 1, + tty: '/dev/pts/1' + }, + capturedAgeMs: 0 + } + ]) + const reader = createWslRelayIdentityReader({ readProcessIdentity }) + const anchor = { + distro: 'Ubuntu', + bootId: '11111111-1111-1111-1111-111111111111', + shellPid: 1, + shellStartTime: 1, + tty: '/dev/pts/1' + } + + await reader.readBatch('Ubuntu', [anchor], { stableUntilReset: true }) + vi.setSystemTime(60_000) + await reader.readBatch('Ubuntu', [anchor]) + vi.setSystemTime(120_000) + await reader.readBatch('Ubuntu', [anchor], { stableUntilReset: true }) + expect(readProcessIdentity).toHaveBeenCalledTimes(2) + vi.setSystemTime(180_000) + await reader.readBatch('Ubuntu', [anchor], { stableUntilReset: true }) + expect(readProcessIdentity).toHaveBeenCalledTimes(2) + }) }) diff --git a/src/main/providers/wsl-relay-identity-reader.ts b/src/main/providers/wsl-relay-identity-reader.ts index a2c079fda2a..40bfc803934 100644 --- a/src/main/providers/wsl-relay-identity-reader.ts +++ b/src/main/providers/wsl-relay-identity-reader.ts @@ -16,7 +16,7 @@ export type WslRelayIdentityReader = { readBatch: ( distro: string, anchors: readonly WslShellProcessAnchor[], - options?: { signal?: AbortSignal; timeoutMs?: number } + options?: { signal?: AbortSignal; timeoutMs?: number; stableUntilReset?: boolean } ) => Promise reset: () => void } @@ -26,21 +26,40 @@ export function createWslRelayIdentityReader( ): WslRelayIdentityReader { const cache = new Map< string, - { at: number; anchors: readonly WslShellProcessAnchor[]; results: WslRelayIdentityResult[] } + { + at: number + anchors: readonly WslShellProcessAnchor[] + results: WslRelayIdentityResult[] + stableUntilReset: boolean + } >() const pending = new Map< string, - Promise<{ anchors: readonly WslShellProcessAnchor[]; results: WslRelayIdentityResult[] }> + { + promise: Promise<{ + anchors: readonly WslShellProcessAnchor[] + results: WslRelayIdentityResult[] + }> + stableUntilReset: boolean + resetGeneration: number + } >() + let resetGeneration = 0 const readBatch = async ( distro: string, anchors: readonly WslShellProcessAnchor[], - options?: { signal?: AbortSignal; timeoutMs?: number } + options?: { signal?: AbortSignal; timeoutMs?: number; stableUntilReset?: boolean } ): Promise => { const key = distro.trim().toLowerCase() const now = Date.now() const prior = cache.get(key) - if (prior && now - prior.at < 500) { + if ( + prior && + ((options?.stableUntilReset === true && prior.stableUntilReset) || now - prior.at < 500) + ) { + if (options?.stableUntilReset === true) { + prior.stableUntilReset = true + } const byAnchor = new Map( prior.anchors.map((anchor, index) => [JSON.stringify(anchor), prior.results[index]!]) ) @@ -56,7 +75,10 @@ export function createWslRelayIdentityReader( } const active = pending.get(key) if (active) { - const result = await active + if (options?.stableUntilReset === true) { + active.stableUntilReset = true + } + const result = await active.promise const byAnchor = new Map( result.anchors.map((anchor, index) => [JSON.stringify(anchor), result.results[index]!]) ) @@ -64,16 +86,27 @@ export function createWslRelayIdentityReader( return anchors.map((anchor) => byAnchor.get(JSON.stringify(anchor))!) } } - const request = manager - .readProcessIdentity(distro, anchors, options) - .then((results) => ({ anchors, results })) - pending.set(key, request) + const requestGeneration = resetGeneration + const activeRequest = { + stableUntilReset: options?.stableUntilReset === true, + resetGeneration: requestGeneration, + promise: manager + .readProcessIdentity(distro, anchors, options) + .then((results) => ({ anchors, results })) + } + pending.set(key, activeRequest) try { - const result = await request - cache.set(key, { at: Date.now(), ...result }) + const result = await activeRequest.promise + if (activeRequest.resetGeneration === resetGeneration) { + cache.set(key, { + at: Date.now(), + stableUntilReset: activeRequest.stableUntilReset || prior?.stableUntilReset === true, + ...result + }) + } return result.results } finally { - if (pending.get(key) === request) { + if (pending.get(key) === activeRequest) { pending.delete(key) } } @@ -82,6 +115,7 @@ export function createWslRelayIdentityReader( read: async (distro, anchor, options) => (await readBatch(distro, [anchor], options))[0]!, readBatch, reset: () => { + resetGeneration += 1 cache.clear() pending.clear() } diff --git a/src/main/runtime/orca-runtime-refresh-floating-workspace-pty-liveness.ts b/src/main/runtime/orca-runtime-refresh-floating-workspace-pty-liveness.ts index c943ec96e99..57dfca4e92a 100644 --- a/src/main/runtime/orca-runtime-refresh-floating-workspace-pty-liveness.ts +++ b/src/main/runtime/orca-runtime-refresh-floating-workspace-pty-liveness.ts @@ -67,11 +67,13 @@ export class OrcaRuntimeWithRefreshFloatingWorkspacePtyLiveness extends OrcaRunt const binding = persistedBindingByPtyId.get(ptyId) if (!pty && binding) { // Why: a live daemon PTY restored from disk needs its pane identity before mobile can issue a safe handle. - pty = this.recordPtyWorktree(ptyId, FLOATING_TERMINAL_WORKTREE_ID, { - connected: true, - tabId: binding.tabId, - paneKey: binding.paneKey - }) + pty = this.withPtyLivenessRefreshBookkeeping(() => + this.recordPtyWorktree(ptyId, FLOATING_TERMINAL_WORKTREE_ID, { + connected: true, + tabId: binding.tabId, + paneKey: binding.paneKey + }) + ) } if (pty) { pty.connected = true diff --git a/src/main/runtime/orca-runtime-refresh-pty-worktree-records-from-controller.ts b/src/main/runtime/orca-runtime-refresh-pty-worktree-records-from-controller.ts index 4f4f9bda307..76fe1be5a19 100644 --- a/src/main/runtime/orca-runtime-refresh-pty-worktree-records-from-controller.ts +++ b/src/main/runtime/orca-runtime-refresh-pty-worktree-records-from-controller.ts @@ -9,6 +9,36 @@ export class OrcaRuntimeWithRefreshPtyWorktreeRecordsFromController extends Orca targetWorktreeId: string | null = null, deadline?: number, signal?: AbortSignal + ): Promise | null> { + const refreshKey = targetWorktreeId === null ? 'aggregate' : `target:${targetWorktreeId}` + const pending = this.ptyLivenessRefreshPromises.get(refreshKey) + if (pending) { + return pending + } + const invalidationSequence = this.ptyLivenessInvalidationSequence + const refresh = this.refreshPtyWorktreeRecordsFromControllerUncoalesced( + resolvedWorktrees, + targetWorktreeId, + deadline, + signal, + invalidationSequence + ) + this.ptyLivenessRefreshPromises.set(refreshKey, refresh) + try { + return await refresh + } finally { + if (this.ptyLivenessRefreshPromises.get(refreshKey) === refresh) { + this.ptyLivenessRefreshPromises.delete(refreshKey) + } + } + } + + private async refreshPtyWorktreeRecordsFromControllerUncoalesced( + resolvedWorktrees: ResolvedWorktree[], + targetWorktreeId: string | null, + deadline?: number, + signal?: AbortSignal, + invalidationSequence = this.ptyLivenessInvalidationSequence ): Promise | null> { this.ptyLivenessRefreshInProgress += 1 try { @@ -20,7 +50,11 @@ export class OrcaRuntimeWithRefreshPtyWorktreeRecordsFromController extends Orca false, signal ? { signal } : undefined ) - if (inventory) { + if ( + inventory && + targetWorktreeId === null && + invalidationSequence === this.ptyLivenessInvalidationSequence + ) { this.ptyLivenessRefreshRequired = false } return inventory ? new Set(inventory.livePtyIds) : null diff --git a/src/main/runtime/orca-runtime-refresh-pty-worktree-records-with-controller-inventory.ts b/src/main/runtime/orca-runtime-refresh-pty-worktree-records-with-controller-inventory.ts index c9e4f38700e..731a5e4d6b5 100644 --- a/src/main/runtime/orca-runtime-refresh-pty-worktree-records-with-controller-inventory.ts +++ b/src/main/runtime/orca-runtime-refresh-pty-worktree-records-with-controller-inventory.ts @@ -217,17 +217,19 @@ export class OrcaRuntimeWithRefreshPtyWorktreeRecordsWithControllerInventory ext } this.restoredOrchestrationAuthorityByPtyId.delete(session.id) if (worktreeId) { - const pty = this.recordPtyWorktree(session.id, worktreeId, { - connected: true, - ...(session.incarnationId ? { incarnationId: session.incarnationId } : {}), - agentSessionOwners: session.incarnationId ? (session.agentSessionOwners ?? []) : [], - ...(session.wslDistro !== undefined - ? { isWsl: Boolean(session.wslDistro), wslDistro: session.wslDistro } - : {}), - ...(restoresExactSurface - ? { tabId: persistedSurface.tabId, paneKey: persistedSurface.paneKey } - : {}) - }) + const pty = this.withPtyLivenessRefreshBookkeeping(() => + this.recordPtyWorktree(session.id, worktreeId, { + connected: true, + ...(session.incarnationId ? { incarnationId: session.incarnationId } : {}), + agentSessionOwners: session.incarnationId ? (session.agentSessionOwners ?? []) : [], + ...(session.wslDistro !== undefined + ? { isWsl: Boolean(session.wslDistro), wslDistro: session.wslDistro } + : {}), + ...(restoresExactSurface + ? { tabId: persistedSurface.tabId, paneKey: persistedSurface.paneKey } + : {}) + }) + ) if (restoresExactSurface && controllerIdentity) { this.rememberRestoredOrchestrationAuthority( pty, diff --git a/src/main/runtime/orca-runtime-runtime-id.ts b/src/main/runtime/orca-runtime-runtime-id.ts index 9b97e7c69b2..b4eade18d98 100644 --- a/src/main/runtime/orca-runtime-runtime-id.ts +++ b/src/main/runtime/orca-runtime-runtime-id.ts @@ -342,9 +342,31 @@ export class OrcaRuntimeWithRuntimeId { protected ptyLivenessRefreshInProgress = 0 + // Concurrent catalog requests share one event-invalidated controller census; + // without this, each request observes the stale flag before the first one clears it. + protected ptyLivenessRefreshPromises = new Map | null>>() + + // Monotonic event fence prevents an in-flight census from clearing a newer invalidation. + protected ptyLivenessInvalidationSequence = 0 + + // Reconciliation bookkeeping must not look like a new PTY event. + protected ptyLivenessRefreshBookkeeping = false + protected invalidatePtyLivenessSnapshot(): void { - if (this.ptyLivenessRefreshInProgress === 0) { - this.ptyLivenessRefreshRequired = true + if (this.ptyLivenessRefreshBookkeeping) { + return + } + this.ptyLivenessRefreshRequired = true + this.ptyLivenessInvalidationSequence += 1 + } + + protected withPtyLivenessRefreshBookkeeping(callback: () => T): T { + const prior = this.ptyLivenessRefreshBookkeeping + this.ptyLivenessRefreshBookkeeping = true + try { + return callback() + } finally { + this.ptyLivenessRefreshBookkeeping = prior } } } diff --git a/src/main/runtime/orca-runtime-tests/mobile-summaries-part-04.spec.ts b/src/main/runtime/orca-runtime-tests/mobile-summaries-part-04.spec.ts index bebc15f36ec..b03faa116e5 100644 --- a/src/main/runtime/orca-runtime-tests/mobile-summaries-part-04.spec.ts +++ b/src/main/runtime/orca-runtime-tests/mobile-summaries-part-04.spec.ts @@ -2,6 +2,7 @@ import { describe, expect, it, vi } from 'vitest' import { OrcaRuntimeService, registerSshGitProvider, + resetPlatform, setPlatform, worktreePathComparison } from '../orca-runtime-test-mocks.spec' @@ -11,6 +12,24 @@ import { getWslRelayIdentityRpcCount, resetWslRelayIdentityRpcCount } from '../../agent-hooks/wsl-hook-relay-manager' +import { listLocalPtyProcesses } from '../../providers/local-pty-session-operations' +import { + clearPtyState, + ptyInitialCwd, + ptyProcesses, + ptyWslDistroById, + ptyWslShellAnchors, + ptyWorktreeId +} from '../../providers/local-pty-provider-state' +import { wslRelayIdentityReader } from '../../providers/wsl-relay-identity-reader' + +function deferred(): { promise: Promise; resolve: (value: T) => void } { + let resolve!: (value: T) => void + const promise = new Promise((next) => { + resolve = next + }) + return { promise, resolve } +} describe('OrcaRuntimeService', () => { it('does not read guest inventories on idle worktree poll ticks', async () => { @@ -42,6 +61,138 @@ describe('OrcaRuntimeService', () => { await runtime.getWorktreePs() expect(relayProcessReads).toHaveBeenCalledOnce() }) + + it('coalesces concurrent and repeated worktree polls after one reconciliation', async () => { + resetWslRelayIdentityRpcCount() + const relayInventory = + deferred<{ id: string; cwd: string; title: string; worktreeId: string }[]>() + const relayProcessReads = vi.fn(() => relayInventory.promise) + const runtime = new OrcaRuntimeService(store) + runtime.setPtyController({ + write: () => true, + kill: () => true, + getForegroundProcess: async () => null, + listProcesses: relayProcessReads + }) + syncSinglePty(runtime) + + const concurrentPolls = [ + runtime.getWorktreePs(), + runtime.getWorktreePs(), + runtime.getWorktreePs() + ] + await vi.waitFor(() => expect(relayProcessReads).toHaveBeenCalledOnce()) + relayInventory.resolve([ + { + id: 'pty-1', + cwd: '/tmp/worktree-a', + title: 'shell', + worktreeId: 'repo-1::/tmp/worktree-a' + } + ]) + await Promise.all(concurrentPolls) + + const identityRpcCountAfterFirstReconciliation = getWslRelayIdentityRpcCount() + expect(identityRpcCountAfterFirstReconciliation).toBe(0) + for (let poll = 0; poll < 5; poll += 1) { + await runtime.getWorktreePs() + } + expect(relayProcessReads).toHaveBeenCalledOnce() + expect(getWslRelayIdentityRpcCount()).toBe(identityRpcCountAfterFirstReconciliation) + }) + + it('keeps WSL inventory reads stable across repeated worktree polls', async () => { + setPlatform('win32') + resetWslRelayIdentityRpcCount() + const ptyId = 'wsl-pty-1' + const anchor = { + distro: 'Ubuntu', + bootId: '11111111-1111-1111-1111-111111111111', + shellPid: 1, + shellStartTime: 1, + tty: '/dev/pts/1' + } + const guestIdentityReads = vi.spyOn(wslRelayIdentityReader, 'readBatch').mockResolvedValue([ + { + status: 'unverifiable' as const, + reason: 'relay_unavailable', + capturedAgeMs: 0 + } + ]) + const fakePty = { process: 'bash' } + ptyProcesses.set(ptyId, fakePty as never) + ptyInitialCwd.set(ptyId, '/tmp/worktree-a') + ptyWorktreeId.set(ptyId, 'repo-1::/tmp/worktree-a') + ptyWslDistroById.set(ptyId, 'Ubuntu') + ptyWslShellAnchors.set(ptyId, anchor) + const guestReads = vi.fn() + const runtime = new OrcaRuntimeService(store) + runtime.setPtyController({ + write: () => true, + kill: () => true, + getForegroundProcess: async () => null, + listProcesses: async () => { + guestReads() + return listLocalPtyProcesses() + } + }) + syncSinglePty(runtime) + + try { + await runtime.getWorktreePs() + const identityRpcCountAfterFirstReconciliation = getWslRelayIdentityRpcCount() + const guestReadsAfterFirstReconciliation = guestReads.mock.calls.length + expect(guestIdentityReads).toHaveBeenCalledOnce() + expect(guestReadsAfterFirstReconciliation).toBe(1) + for (let poll = 0; poll < 5; poll += 1) { + await runtime.getWorktreePs() + } + expect(guestReads).toHaveBeenCalledTimes(guestReadsAfterFirstReconciliation) + expect(guestIdentityReads).toHaveBeenCalledOnce() + expect(getWslRelayIdentityRpcCount()).toBe(identityRpcCountAfterFirstReconciliation) + } finally { + guestIdentityReads.mockRestore() + clearPtyState(ptyId) + resetPlatform() + } + }) + + it('does not let targeted liveness satisfy a pending aggregate poll', async () => { + const relayInventory = + deferred<{ id: string; cwd: string; title: string; worktreeId: string }[]>() + const relayProcessReads = vi.fn(() => relayInventory.promise) + const runtime = new OrcaRuntimeService(store) + runtime.setPtyController({ + write: () => true, + kill: () => true, + getForegroundProcess: async () => null, + listProcesses: relayProcessReads + }) + syncSinglePty(runtime) + const refresh = ( + runtime as unknown as { + refreshPtyWorktreeRecordsFromController: ( + worktrees: readonly unknown[], + targetWorktreeId: string + ) => Promise | null> + } + ).refreshPtyWorktreeRecordsFromController + + const targeted = refresh.call(runtime, [], 'repo-1::/tmp/worktree-a') + await vi.waitFor(() => expect(relayProcessReads).toHaveBeenCalledOnce()) + relayInventory.resolve([ + { + id: 'pty-1', + cwd: '/tmp/worktree-a', + title: 'shell', + worktreeId: 'repo-1::/tmp/worktree-a' + } + ]) + await targeted + + await runtime.getWorktreePs() + expect(relayProcessReads).toHaveBeenCalledTimes(2) + }) it('keeps pinned and unread worktrees when active rows fill the mobile summary limit', async () => { setPlatform('win32') const remoteRepo = {