From fe785ceabc012b22803838f28cd609d2bf268876 Mon Sep 17 00:00:00 2001 From: Merge Sim Date: Tue, 1 Sep 2026 17:44:28 -0700 Subject: [PATCH] perf(wsl): bound guest inventory refresh and correlation --- src/main/daemon/daemon-pty-router.ts | 5 +- .../daemon/daemon-pty-session-inventory.ts | 41 +++++++-- .../daemon/degraded-daemon-pty-provider.ts | 5 +- src/main/ipc/pty/runtime/operations.ts | 8 +- .../providers/local-pty-provider-state.ts | 2 + src/main/providers/local-pty-provider.ts | 4 +- .../providers/local-pty-session-operations.ts | 39 ++++++-- src/main/providers/pty-provider-contract.ts | 2 +- src/main/providers/ssh-pty-provider.ts | 20 ++++- ...wsl-guest-foreground-process-resolution.ts | 63 +++++++++---- .../wsl-guest-process-inventory.test.ts | 74 +++++++++++++++ .../providers/wsl-guest-process-inventory.ts | 88 +++++++++++++++--- src/main/runtime/orca-runtime.test.ts | 25 ++++++ src/main/runtime/orca-runtime.ts | 89 ++++++++++++++++--- .../rpc/methods/worktree-catalog-methods.ts | 14 +-- 15 files changed, 407 insertions(+), 72 deletions(-) diff --git a/src/main/daemon/daemon-pty-router.ts b/src/main/daemon/daemon-pty-router.ts index ab30878c9d8..0a93b417755 100644 --- a/src/main/daemon/daemon-pty-router.ts +++ b/src/main/daemon/daemon-pty-router.ts @@ -197,7 +197,10 @@ export class DaemonPtyRouter implements IPtyProvider { await this.current.revive(state) } - async listProcesses(opts?: { deadlineMs?: number }): Promise { + async listProcesses(opts?: { + deadlineMs?: number + signal?: AbortSignal + }): Promise { // Why: runtime exact-stop/liveness flows must fail closed if any adapter // cannot provide a trustworthy process list. const results = await Promise.all( diff --git a/src/main/daemon/daemon-pty-session-inventory.ts b/src/main/daemon/daemon-pty-session-inventory.ts index 8735a2efe67..22b0513bc23 100644 --- a/src/main/daemon/daemon-pty-session-inventory.ts +++ b/src/main/daemon/daemon-pty-session-inventory.ts @@ -16,13 +16,18 @@ import { PtyProcessListAdmission } from '../providers/pty-process-list-admission import type { PtyProcessInfo } from '../providers/types' import type { ForegroundProcessEvidence } from '../../shared/foreground-process-evidence' import { + createWslGuestProcessIndexes, readWslGuestProcessInventory, resolveWslGuestForegroundProcess, + WSL_GUEST_INVENTORY_MAX_CONCURRENCY, type WslGuestProcessInventoryRead } from '../providers/wsl-guest-process-inventory' export abstract class DaemonPtySessionInventory extends DaemonPtyProcessInspection { - async listProcesses(opts?: { deadlineMs?: number }): Promise { + async listProcesses(opts?: { + deadlineMs?: number + signal?: AbortSignal + }): Promise { // Why: snapshotted before the request so ids spawned mid-flight can never // be reconciled away below. const preRequestActiveIds = new Set(this.activeSessionIds) @@ -60,12 +65,33 @@ export abstract class DaemonPtySessionInventory extends DaemonPtyProcessInspecti .filter((session) => session.wslShellAnchor) .map((session) => session.wslDistro as string) ) + const distroList = [...distros] + let nextDistroIndex = 0 + const readNextDistro = async (): Promise => { + while (nextDistroIndex < distroList.length) { + const distro = distroList[nextDistroIndex++]! + wslByDistro.set( + distro, + await readWslGuestProcessInventory(distro, { + deadlineMs: opts?.deadlineMs, + signal: opts?.signal + }) + ) + } + } await Promise.all( - [...distros].map(async (distro) => { - wslByDistro.set(distro, await readWslGuestProcessInventory(distro)) - }) + Array.from( + { length: Math.min(WSL_GUEST_INVENTORY_MAX_CONCURRENCY, distroList.length) }, + () => readNextDistro() + ) ) } + const indexesByDistro = new Map>() + for (const [distro, inventory] of wslByDistro) { + if (inventory.status === 'ok') { + indexesByDistro.set(distro, createWslGuestProcessIndexes(inventory.inventory)) + } + } for (const session of result.sessions) { if (!session.isAlive) { continue @@ -75,7 +101,12 @@ export abstract class DaemonPtySessionInventory extends DaemonPtyProcessInspecti const inventory = session.wslDistro ? wslByDistro.get(session.wslDistro) : undefined const resolution = session.wslDistro && session.wslShellAnchor && inventory?.status === 'ok' - ? resolveWslGuestForegroundProcess(inventory.inventory, session.wslShellAnchor) + ? resolveWslGuestForegroundProcess( + inventory.inventory, + session.wslShellAnchor, + indexesByDistro.get(session.wslDistro) ?? + createWslGuestProcessIndexes(inventory.inventory) + ) : null const foregroundProcessEvidence: ForegroundProcessEvidence | undefined = session.wslDistro ? resolution?.status === 'live' diff --git a/src/main/daemon/degraded-daemon-pty-provider.ts b/src/main/daemon/degraded-daemon-pty-provider.ts index f0e183bcb20..1d8bd439936 100644 --- a/src/main/daemon/degraded-daemon-pty-provider.ts +++ b/src/main/daemon/degraded-daemon-pty-provider.ts @@ -206,7 +206,10 @@ export class DegradedDaemonPtyProvider implements IPtyProvider { await this.fallback.revive(state) } - async listProcesses(opts?: { deadlineMs?: number }): Promise { + async listProcesses(opts?: { + deadlineMs?: number + signal?: AbortSignal + }): Promise { const results = await Promise.all( this.allProviders().map((provider) => provider.listProcesses(opts)) ) diff --git a/src/main/ipc/pty/runtime/operations.ts b/src/main/ipc/pty/runtime/operations.ts index 00db3380309..b139e737930 100644 --- a/src/main/ipc/pty/runtime/operations.ts +++ b/src/main/ipc/pty/runtime/operations.ts @@ -230,7 +230,7 @@ function markSshInventoryUnverifiable( export async function listProcessesWithHostScopeFromRuntimeController( deps: PtyRuntimeControllerDeps, - opts?: { deadlineMs?: number } + opts?: { deadlineMs?: number; signal?: AbortSignal } ): Promise<{ processes: PtyProcessInfo[]; hostIds: ExecutionHostId[] }> { const providerSessions = await Promise.all( registeredPtyProviders().map(async ({ provider, connectionId }) => { @@ -239,7 +239,7 @@ export async function listProcessesWithHostScopeFromRuntimeController( : LOCAL_EXECUTION_HOST_ID try { return { - processes: await (connectionId ? provider.listProcesses(opts) : provider.listProcesses()), + processes: await provider.listProcesses(opts), hostId } } catch (error) { @@ -261,10 +261,10 @@ export async function listProcessesWithHostScopeFromRuntimeController( export async function listProcessesFromRuntimeController( deps: PtyRuntimeControllerDeps, connectionId?: string | null, - opts?: { deadlineMs?: number } + opts?: { deadlineMs?: number; signal?: AbortSignal } ) { if (connectionId === null) { - return localProvider.listProcesses() + return localProvider.listProcesses(opts) } if (connectionId !== undefined) { try { diff --git a/src/main/providers/local-pty-provider-state.ts b/src/main/providers/local-pty-provider-state.ts index a63c5e993f7..87a0b2dfdac 100644 --- a/src/main/providers/local-pty-provider-state.ts +++ b/src/main/providers/local-pty-provider-state.ts @@ -4,6 +4,7 @@ import type { PtyStartupIngress } from '../../shared/pty-startup-ingress' import type { TerminalExitCause } from '../../shared/terminal-exit-cause' import { normalizeLocalCallerSessionId } from './local-pty-launch-helpers' import type { WslShellProcessAnchor } from '../../shared/wsl-shell-process-anchor' +import { resetWslGuestProcessInventory } from './wsl-guest-process-inventory' export type PtyShutdownOperation = { promise: Promise @@ -132,6 +133,7 @@ export function clearPtyState(id: string): void { ptyTerminationMode.delete(id) ptyReportsChildExitStatus.delete(id) ptyPhysicalExits.delete(id) + resetWslGuestProcessInventory() } /** diff --git a/src/main/providers/local-pty-provider.ts b/src/main/providers/local-pty-provider.ts index d08e437add3..74b45cadf3c 100644 --- a/src/main/providers/local-pty-provider.ts +++ b/src/main/providers/local-pty-provider.ts @@ -136,8 +136,8 @@ export class LocalPtyProvider implements IPtyProvider { /* re-spawning handles local revival */ } - listProcesses(): Promise { - return listLocalPtyProcesses() + listProcesses(opts?: { deadlineMs?: number; signal?: AbortSignal }): Promise { + return listLocalPtyProcesses(opts) } getDefaultShell(): Promise { diff --git a/src/main/providers/local-pty-session-operations.ts b/src/main/providers/local-pty-session-operations.ts index 748b9334a0e..80c8b4b0086 100644 --- a/src/main/providers/local-pty-session-operations.ts +++ b/src/main/providers/local-pty-session-operations.ts @@ -24,8 +24,10 @@ import { import type { LocalPtyProviderOptions } from './local-pty-provider-types' import type { PtyProcessInfo } from './types' import { + createWslGuestProcessIndexes, readWslGuestProcessInventory, resolveWslGuestForegroundProcess, + WSL_GUEST_INVENTORY_MAX_CONCURRENCY, type WslGuestProcessInventoryRead } from './wsl-guest-process-inventory' import type { ForegroundProcessEvidence } from '../../shared/foreground-process-evidence' @@ -122,7 +124,10 @@ export function closeLocalPtyStartupQueryAuthority(id: string): number { return startupIngressByPty.get(id)?.closeQueryAuthority() ?? 0 } -export async function listLocalPtyProcesses(): Promise { +export async function listLocalPtyProcesses(opts?: { + deadlineMs?: number + signal?: AbortSignal +}): Promise { const entries = Array.from(ptyProcesses.entries()) const evidenceEpoch = Date.now() const wslByDistro = new Map() @@ -137,11 +142,31 @@ export async function listLocalPtyProcesses(): Promise { } } const inventories = new Map() + const distros = [...wslByDistro.keys()] + let nextDistroIndex = 0 + const readNextDistro = async (): Promise => { + while (nextDistroIndex < distros.length) { + const distro = distros[nextDistroIndex++]! + inventories.set( + distro, + await readWslGuestProcessInventory(distro, { + deadlineMs: opts?.deadlineMs, + signal: opts?.signal + }) + ) + } + } await Promise.all( - [...wslByDistro.keys()].map(async (distro) => { - inventories.set(distro, await readWslGuestProcessInventory(distro)) - }) + Array.from({ length: Math.min(WSL_GUEST_INVENTORY_MAX_CONCURRENCY, distros.length) }, () => + readNextDistro() + ) ) + const indexesByDistro = new Map>() + for (const [distro, read] of inventories) { + if (read.status === 'ok') { + indexesByDistro.set(distro, createWslGuestProcessIndexes(read.inventory)) + } + } return entries.flatMap(([id, proc]) => { // Inventory reads are asynchronous; a PTY may have exited while they ran. @@ -157,7 +182,11 @@ export async function listLocalPtyProcesses(): Promise { const anchor = ptyWslShellAnchors.get(id) const resolution = read?.status === 'ok' && anchor - ? resolveWslGuestForegroundProcess(read.inventory, anchor) + ? resolveWslGuestForegroundProcess( + read.inventory, + anchor, + indexesByDistro.get(distro) ?? createWslGuestProcessIndexes(read.inventory) + ) : { status: 'unverifiable' as const, reason: read?.status === 'unverifiable' ? read.reason : 'anchor_missing' diff --git a/src/main/providers/pty-provider-contract.ts b/src/main/providers/pty-provider-contract.ts index 35fca5b0b34..d17a68411af 100644 --- a/src/main/providers/pty-provider-contract.ts +++ b/src/main/providers/pty-provider-contract.ts @@ -222,7 +222,7 @@ export type IPtyProvider = { serialize(ids: string[]): Promise revive(state: string): Promise // Why: deadlineMs bounds the underlying RPC exactly like shutdown's deadlineMs. - listProcesses(opts?: { deadlineMs?: number }): Promise + listProcesses(opts?: { deadlineMs?: number; signal?: AbortSignal }): Promise getDefaultShell(): Promise getProfiles(): Promise<{ name: string; path: string }[]> onData(callback: (payload: PtyDataEvent) => void): () => void diff --git a/src/main/providers/ssh-pty-provider.ts b/src/main/providers/ssh-pty-provider.ts index a8e4218ce63..32a467d7a32 100644 --- a/src/main/providers/ssh-pty-provider.ts +++ b/src/main/providers/ssh-pty-provider.ts @@ -25,8 +25,17 @@ import type { PtyProcessInspection } from './pty-process-inspection' import { writeToSshPty, writeToSshPtyWithSettlement } from './ssh-pty-write' // Why: sequential relay teardown calls share one absolute budget; convert to the mux-relative timeout only at dispatch. -function relayTimeoutOptions(deadlineMs: number | undefined): { timeoutMs: number } | undefined { - return deadlineMs === undefined ? undefined : { timeoutMs: Math.max(1, deadlineMs - Date.now()) } +function relayTimeoutOptions( + deadlineMs: number | undefined, + signal?: AbortSignal +): { timeoutMs?: number; signal?: AbortSignal } | undefined { + if (deadlineMs === undefined && signal === undefined) { + return undefined + } + return { + ...(deadlineMs === undefined ? {} : { timeoutMs: Math.max(1, deadlineMs - Date.now()) }), + ...(signal ? { signal } : {}) + } } /** Remote PTY provider that proxies IPtyProvider operations through the relay. */ @@ -277,11 +286,14 @@ export class SshPtyProvider implements IPtyProvider { await this.mux.request('pty.revive', { state }) } - async listProcesses(opts?: { deadlineMs?: number }): Promise { + async listProcesses(opts?: { + deadlineMs?: number + signal?: AbortSignal + }): Promise { const result = await this.mux.request( 'pty.listProcesses', undefined, - relayTimeoutOptions(opts?.deadlineMs) + relayTimeoutOptions(opts?.deadlineMs, opts?.signal) ) const processes = mapSshPtyProcessList(result as PtyProcessInfo[], (id) => this.toAppPtyId(id)) for (const process of processes) { diff --git a/src/main/providers/wsl-guest-foreground-process-resolution.ts b/src/main/providers/wsl-guest-foreground-process-resolution.ts index e8959cdd618..f76544aa1eb 100644 --- a/src/main/providers/wsl-guest-foreground-process-resolution.ts +++ b/src/main/providers/wsl-guest-foreground-process-resolution.ts @@ -11,14 +11,54 @@ export type WslGuestForegroundResolution = | { status: 'live'; processName: string | null; anchor: WslGuestProcessAnchor } | { status: 'unverifiable'; reason: string } +export type WslGuestProcessIndexes = { + byPid: ReadonlyMap + byForegroundGroup: ReadonlyMap + multiplexerRows: readonly WslGuestProcessRow[] +} + function normalizeTty(tty: string): string { return tty.startsWith('/dev/') ? tty : tty === '?' ? '' : `/dev/${tty}` } +function foregroundGroupKey(pgid: number, tty: string): string { + return `${pgid}\u0000${normalizeTty(tty)}` +} + +const isMultiplexerCommand = (command: string): boolean => + /(?:^|\s)(?:tmux|screen)(?:\s|$)/.test(command) + +/** Build the indexes shared by every pane resolution for one inventory. */ +export function createWslGuestProcessIndexes( + inventory: WslGuestProcessInventory +): WslGuestProcessIndexes { + const byPid = new Map() + const groups = new Map() + const multiplexerRows: WslGuestProcessRow[] = [] + for (const row of inventory.rows) { + // Preserve the resolver's historical `rows.find(pid)` first-match rule. + if (!byPid.has(row.pid)) { + byPid.set(row.pid, row) + } + const key = foregroundGroupKey(row.pgid, row.tty) + const group = groups.get(key) + if (group) { + group.push(row) + } else { + groups.set(key, [row]) + } + if (isMultiplexerCommand(row.command)) { + multiplexerRows.push(row) + } + } + return { byPid, byForegroundGroup: groups, multiplexerRows } +} + /** Correlate one shell anchor to its foreground group and strict agent recognizer. */ export function resolveWslGuestForegroundProcess( inventory: WslGuestProcessInventory, - anchor: WslGuestProcessAnchor + anchor: WslGuestProcessAnchor, + indexes: WslGuestProcessIndexes = createWslGuestProcessIndexes(inventory) ): WslGuestForegroundResolution { if (inventory.distro.toLowerCase() !== anchor.distro.toLowerCase()) { return { status: 'unverifiable', reason: 'distro_mismatch' } @@ -26,7 +66,7 @@ export function resolveWslGuestForegroundProcess( if (inventory.bootId !== anchor.bootId) { return { status: 'unverifiable', reason: 'boot_id_mismatch' } } - const shell = inventory.rows.find((row) => row.pid === anchor.shellPid) + const shell = indexes.byPid.get(anchor.shellPid) if (!shell) { return { status: 'unverifiable', reason: 'anchor_missing' } } @@ -40,20 +80,15 @@ export function resolveWslGuestForegroundProcess( if (shell.tpgid <= 0) { return { status: 'unverifiable', reason: 'foreground_group_missing' } } - const group = inventory.rows.filter( - (row) => row.pgid === shell.tpgid && normalizeTty(row.tty) === tty - ) + const group = indexes.byForegroundGroup.get(foregroundGroupKey(shell.tpgid, tty)) ?? [] if (group.length === 0) { return { status: 'unverifiable', reason: 'foreground_group_missing' } } // Multiplexers move the real command to another PTY/session. Without a // session-aware anchor, the outer shell cannot make a truthful claim. - const isMultiplexer = (command: string): boolean => - /(?:^|\s)(?:tmux|screen)(?:\s|$)/.test(command) - if (group.some((row) => isMultiplexer(row.command))) { + if (group.some((row) => isMultiplexerCommand(row.command))) { return { status: 'unverifiable', reason: 'multiplexer_boundary' } } - const byPid = new Map(inventory.rows.map((row) => [row.pid, row])) const isShellDescendant = (row: WslGuestProcessRow): boolean => { const seen = new Set() let current: WslGuestProcessRow | undefined = row @@ -62,17 +97,13 @@ export function resolveWslGuestForegroundProcess( return true } seen.add(current.pid) - current = byPid.get(current.ppid) + current = indexes.byPid.get(current.ppid) } return false } if ( - inventory.rows.some( - (row) => - row.pid !== shell.pid && - normalizeTty(row.tty) !== tty && - isMultiplexer(row.command) && - isShellDescendant(row) + indexes.multiplexerRows.some( + (row) => row.pid !== shell.pid && normalizeTty(row.tty) !== tty && isShellDescendant(row) ) ) { return { status: 'unverifiable', reason: 'multiplexer_boundary' } diff --git a/src/main/providers/wsl-guest-process-inventory.test.ts b/src/main/providers/wsl-guest-process-inventory.test.ts index ac4a4ae1c3c..f797b066099 100644 --- a/src/main/providers/wsl-guest-process-inventory.test.ts +++ b/src/main/providers/wsl-guest-process-inventory.test.ts @@ -166,6 +166,31 @@ describe('WSL guest process inventory', () => { ).toEqual({ status: 'unverifiable', reason: 'pid_reused' }) }) + it('resolves against a prebuilt inventory index without rescanning rows', () => { + const inventory = parseWslGuestProcessInventoryPayload( + payload( + [ + 'row 100 90 90 100 100 pts/0 Ss+ 12345 bash', + 'row 101 100 100 101 101 pts/0 Sl+ 54321 codex' + ].join('\n'), + 2 + ), + 'Ubuntu' + ) + const indexes = { + byPid: new Map(), + byForegroundGroup: new Map(), + multiplexerRows: [] + } + expect( + resolveWslGuestForegroundProcess( + inventory, + { distro: 'Ubuntu', bootId, shellPid: 100, shellStartTime: 12345, tty: '/dev/pts/0' }, + indexes + ) + ).toEqual({ status: 'unverifiable', reason: 'anchor_missing' }) + }) + it('does not claim identity across a multiplexer boundary', () => { const inventory = parseWslGuestProcessInventoryPayload( payload( @@ -212,6 +237,55 @@ describe('WSL guest process inventory', () => { expect(calls).toBe(3) }) + it('bounds the derived inventory cache and evicts the least-recently-used distro', async () => { + let calls = 0 + const reader = createWslGuestProcessInventoryReader({ + run: async (distro) => { + calls += 1 + return { status: 'ok', inventory: { distro, bootId, rows: [] } } + } + }) + for (let index = 0; index < 40; index += 1) { + await reader.read(`distro-${index}`) + } + expect(calls).toBe(40) + await reader.read('distro-0') + expect(calls).toBe(41) + await reader.read('distro-39') + expect(calls).toBe(41) + }) + + it('passes caller cancellation and deadline through to the guest probe', async () => { + let observedOpts: { deadlineMs?: number; signal?: AbortSignal } | undefined + const run = vi.fn( + async (distro: string, opts?: { deadlineMs?: number; signal?: AbortSignal }) => { + observedOpts = opts + return { status: 'ok' as const, inventory: { distro, bootId, rows: [] } } + } + ) + const reader = createWslGuestProcessInventoryReader({ run, now: () => 0 }) + const signal = new AbortController().signal + await reader.read('Ubuntu', { deadlineMs: 1234, signal }) + expect(run).toHaveBeenCalledOnce() + expect(observedOpts?.deadlineMs).toBe(1234) + expect(observedOpts?.signal).toBe(signal) + }) + + it('bounds the guest process timeout by the caller deadline', async () => { + runProcessMock.mockResolvedValue({ code: 127, stdout: '', stderr: '', timedOut: false }) + resetWslGuestProcessInventoryForTests() + const signal = new AbortController().signal + const deadlineMs = Date.now() + 1_000 + await readWslGuestProcessInventory('Ubuntu', { deadlineMs, signal }) + const spec = runProcessMock.mock.calls[0]?.[0] as { + timeoutMs?: number + signal?: AbortSignal + } + expect(spec.timeoutMs).toBeGreaterThan(0) + expect(spec.timeoutMs).toBeLessThanOrEqual(1_000) + expect(spec.signal).toBe(signal) + }) + it.each([1, 8, 32])('uses one guest inventory for a %s-pane burst', async (paneCount) => { let calls = 0 const reader = createWslGuestProcessInventoryReader({ diff --git a/src/main/providers/wsl-guest-process-inventory.ts b/src/main/providers/wsl-guest-process-inventory.ts index 06635872ff7..1dac6261bf0 100644 --- a/src/main/providers/wsl-guest-process-inventory.ts +++ b/src/main/providers/wsl-guest-process-inventory.ts @@ -12,10 +12,14 @@ export type { WslGuestProcessInventory, WslGuestProcessRow } from './wsl-guest-process-inventory-parser' -export { resolveWslGuestForegroundProcess } from './wsl-guest-foreground-process-resolution' +export { + createWslGuestProcessIndexes, + resolveWslGuestForegroundProcess +} from './wsl-guest-foreground-process-resolution' export type { WslGuestForegroundResolution, - WslGuestProcessAnchor + WslGuestProcessAnchor, + WslGuestProcessIndexes } from './wsl-guest-foreground-process-resolution' export type WslGuestProcessInventoryRead = @@ -33,6 +37,8 @@ export type WslGuestProcessInventoryFailureReason = const INVENTORY_TIMEOUT_MS = 5_000 const INVENTORY_MAX_OUTPUT_BYTES = 4 * 1024 * 1024 const INVENTORY_TTL_MS = 500 +export const WSL_GUEST_INVENTORY_MAX_CONCURRENCY = 4 +const INVENTORY_CACHE_MAX_DISTROS = 32 const CAPTURE_NONCE_ENV = 'ORCA_WSL_CAPTURE_NONCE' /** @@ -83,40 +89,82 @@ export const WSL_GUEST_INVENTORY_SCRIPT = [ ].join('\n') type ReaderDeps = { - run?: (distro: string) => Promise + run?: ( + distro: string, + opts?: { deadlineMs?: number; signal?: AbortSignal } + ) => Promise now?: () => number ttlMs?: number } +export type WslGuestProcessInventoryReadOptions = { + deadlineMs?: number + signal?: AbortSignal +} + /** Construct a per-distro single-flight/TTL reader; exported for deterministic tests. */ export function createWslGuestProcessInventoryReader(deps: ReaderDeps = {}): { - read: (distro: string) => Promise + read: ( + distro: string, + opts?: WslGuestProcessInventoryReadOptions + ) => Promise reset: () => void } { const now = deps.now ?? (() => Date.now()) const ttlMs = deps.ttlMs ?? INVENTORY_TTL_MS const cached = new Map() const inFlight = new Map>() + let resetGeneration = 0 const run = deps.run ?? runWslGuestProcessInventory - const read = (distro: string): Promise => { + const read = ( + distro: string, + opts?: WslGuestProcessInventoryReadOptions + ): Promise => { const cleanedDistro = distro.trim() const key = cleanedDistro.toLowerCase() + const currentTime = now() + for (const [cachedKey, entry] of cached) { + if (currentTime - entry.at >= ttlMs) { + cached.delete(cachedKey) + } + } const prior = cached.get(key) - if (prior && now() - prior.at < ttlMs) { + if (prior) { + // Touch the entry so the map order is a true LRU order while retaining + // the completion timestamp used by the TTL. + cached.delete(key) + cached.set(key, prior) return Promise.resolve(prior.value) } + if ( + opts?.signal?.aborted || + (opts?.deadlineMs !== undefined && opts.deadlineMs <= currentTime) + ) { + return Promise.resolve({ status: 'unverifiable', reason: 'capture_timed_out' }) + } const active = inFlight.get(key) if (active) { return active } - const pending = run(cleanedDistro) + const generationAtStart = resetGeneration + const pending = run(cleanedDistro, opts) .catch((): WslGuestProcessInventoryRead => ({ status: 'unverifiable', reason: 'capture_failed' })) .then((value) => { + if (generationAtStart !== resetGeneration) { + return value + } cached.set(key, { value, at: now() }) + while (cached.size > INVENTORY_CACHE_MAX_DISTROS) { + const oldest = cached.keys().next().value + if (oldest === undefined) { + break + } + cached.delete(oldest) + } return value }) .finally(() => { @@ -130,13 +178,20 @@ export function createWslGuestProcessInventoryReader(deps: ReaderDeps = {}): { return { read, reset: () => { + resetGeneration += 1 cached.clear() inFlight.clear() } } } -async function runWslGuestProcessInventory(distro: string): Promise { +async function runWslGuestProcessInventory( + distro: string, + opts?: WslGuestProcessInventoryReadOptions +): Promise { + if (opts?.signal?.aborted || (opts?.deadlineMs !== undefined && opts.deadlineMs <= Date.now())) { + return { status: 'unverifiable', reason: 'capture_timed_out' } + } const captureNonce = `${Date.now().toString(36)}${Math.random().toString(36).slice(2, 10)}` const captured = buildWslCapturedLoginShellCommand(WSL_GUEST_INVENTORY_SCRIPT, captureNonce, { nonceEnvVar: CAPTURE_NONCE_ENV @@ -162,13 +217,17 @@ async function runWslGuestProcessInventory(distro: string): Promise { - return defaultReader.read(distro) + return defaultReader.read(distro, opts) } -export function resetWslGuestProcessInventoryForTests(): void { +export function resetWslGuestProcessInventory(): void { defaultReader.reset() } + +export const resetWslGuestProcessInventoryForTests = resetWslGuestProcessInventory diff --git a/src/main/runtime/orca-runtime.test.ts b/src/main/runtime/orca-runtime.test.ts index e4c39a32583..f9621368534 100644 --- a/src/main/runtime/orca-runtime.test.ts +++ b/src/main/runtime/orca-runtime.test.ts @@ -38157,6 +38157,31 @@ describe('OrcaRuntimeService', () => { }) }) + it('does not read guest inventories on idle worktree poll ticks', async () => { + const guestInventoryReads = vi.fn() + const runtime = new OrcaRuntimeService(store) + runtime.setPtyController({ + write: () => true, + kill: () => true, + getForegroundProcess: async () => null, + listProcesses: async () => { + guestInventoryReads() + return [] + } + }) + + await runtime.getWorktreePs() + await runtime.getWorktreePs() + await runtime.getWorktreePs() + expect(guestInventoryReads).not.toHaveBeenCalled() + + runtime.notifyBranchRenamed(TEST_REPO_ID) + await runtime.getWorktreePs() + expect(guestInventoryReads).toHaveBeenCalledOnce() + await runtime.getWorktreePs() + expect(guestInventoryReads).toHaveBeenCalledOnce() + }) + it('reads the linked-PR state from the renderer repoId-keyed GitHub cache', async () => { // Regression: renderer keys the PR cache by repoId::branch; reading by path::branch missed every entry (muted mobile badge). const runtimeStore = { diff --git a/src/main/runtime/orca-runtime.ts b/src/main/runtime/orca-runtime.ts index c8a9c7eb376..3c6670d22ae 100644 --- a/src/main/runtime/orca-runtime.ts +++ b/src/main/runtime/orca-runtime.ts @@ -2160,9 +2160,9 @@ type RuntimePtyController = { // the mux's own 30s default and blows every inventory refresh (STA-517). listProcesses?( connectionId?: string | null, - opts?: { deadlineMs?: number } + opts?: { deadlineMs?: number; signal?: AbortSignal } ): Promise - listProcessesWithHostScope?(opts?: { deadlineMs?: number }): Promise<{ + listProcessesWithHostScope?(opts?: { deadlineMs?: number; signal?: AbortSignal }): Promise<{ processes: PtyProcessInfo[] hostIds: ExecutionHostId[] }> @@ -3566,6 +3566,16 @@ export class OrcaRuntimeService { // record so a close/stop receipt can still say the stop was unconfirmed. private ptyLivenessVerdictByPtyId = new Map() private ptyLivenessObservationSequence = 0 + // Catalog polls reuse the last controller census until a PTY lifecycle or + // output event invalidates it; this keeps idle mobile polls read-free. + private ptyLivenessRefreshRequired = false + private ptyLivenessRefreshInProgress = 0 + + private invalidatePtyLivenessSnapshot(): void { + if (this.ptyLivenessRefreshInProgress === 0) { + this.ptyLivenessRefreshRequired = true + } + } private readonly pairedRendererSessionOwnedPtyIds = new Set() private wslDistroByPtyId = new Map() private titleObservationSequence = 0 @@ -6705,6 +6715,28 @@ export class OrcaRuntimeService { // instead of tunneling back through renderer IPC, or live handles could // drift from the process they are supposed to control during reloads. this.ptyController = controller + // A controller attached after restart must reconcile persisted PTY ids once; + // an otherwise idle runtime with no persisted terminals stays read-free. + if (controller && this.hasPersistedPtyReferences()) { + this.invalidatePtyLivenessSnapshot() + } + } + + private hasPersistedPtyReferences(): boolean { + const session = this.store?.getWorkspaceSession?.() + if (!session) { + return false + } + if ( + Object.values(session.tabsByWorktree ?? {}).some((tabs) => + tabs.some((tab) => tab.ptyId !== null) + ) + ) { + return true + } + return Object.values(session.terminalLayoutsByTabId ?? {}).some((layout) => + Object.values(layout?.ptyIdsByLeafId ?? {}).some((ptyId) => Boolean(ptyId)) + ) } setNotifier(notifier: RuntimeNotifier | null): void { @@ -6965,6 +6997,7 @@ export class OrcaRuntimeService { } private notifyWorktreesChanged(repoId: string): void { + this.invalidatePtyLivenessSnapshot() this.notifier?.worktreesChanged(repoId) this.emitClientEvent({ type: 'worktreesChanged', repoId }) } @@ -6992,6 +7025,7 @@ export class OrcaRuntimeService { } private notifyReposChanged(): void { + this.invalidatePtyLivenessSnapshot() wakeFolderRepoGitUpgradeWatch() this.notifier?.reposChanged() this.emitClientEvent({ type: 'reposChanged' }) @@ -13455,6 +13489,7 @@ export class OrcaRuntimeService { captureModelReceipt?: (completion: Promise) => void, sourceRanges?: readonly TerminalOutputSourceRange[] ): number { + this.invalidatePtyLivenessSnapshot() const outputSequence = (this.ptyOutputSequenceById.get(ptyId) ?? 0) + sequenceChars this.ptyOutputSequenceById.set(ptyId, outputSequence) this.providerModeTrackersByPtyId.get(ptyId)?.scan(data) @@ -17730,6 +17765,7 @@ export class OrcaRuntimeService { if (exitIncarnationId && pty?.incarnationId && exitIncarnationId !== pty.incarnationId) { return } + this.invalidatePtyLivenessSnapshot() // Why intent first: a requested stop can still be delivered by the provider's // own exit event, whose status looks exactly like a natural finish. // @@ -22619,7 +22655,8 @@ export class OrcaRuntimeService { async getWorktreePs( limit = DEFAULT_WORKTREE_PS_LIMIT, - sourceDefaultsSupported = true + sourceDefaultsSupported = true, + opts?: { deadlineMs?: number; signal?: AbortSignal } ): Promise<{ worktrees: RuntimeWorktreePsSummary[] totalCount: number @@ -22651,7 +22688,14 @@ export class OrcaRuntimeService { ) // Why: worktree.ps backs the mobile sidebar, so it must use the same // host-owned imported-worktree visibility gate as worktree.list/desktop. - const freshPtyLiveness = await this.refreshPtyWorktreeRecordsFromController(resolvedWorktrees) + const freshPtyLiveness = this.ptyLivenessRefreshRequired + ? await this.refreshPtyWorktreeRecordsFromController( + resolvedWorktrees, + null, + opts?.deadlineMs, + opts?.signal + ) + : null const repoById = new Map((this.store?.getRepos() ?? []).map((repo) => [repo.id, repo])) const platformByRepoId = resolvedWorktreeSnapshot.platformByRepoId const summaries = new Map() @@ -35344,6 +35388,7 @@ export class OrcaRuntimeService { this.clientSessionTabSelections.migrateWorktree(oldWorktreeId, newWorktreeId) this.invalidateResolvedWorktreeCache() this.invalidateWorktreeScanCacheForRepo(repoId) + this.invalidatePtyLivenessSnapshot() this.notifier?.worktreesChanged(repoId, { oldWorktreeId, newWorktreeId }) // Mirror notifyBranchRenamed so in-process onClientEvent listeners also see the rename. this.emitClientEvent({ type: 'worktreesChanged', repoId }) @@ -35375,6 +35420,7 @@ export class OrcaRuntimeService { > > = {} ): RuntimePtyWorktreeRecord { + this.invalidatePtyLivenessSnapshot() let pty = this.ptysById.get(ptyId) if (!pty) { const titleObservedAt = state.title ? this.nextTitleObservationSequence() : null @@ -35537,14 +35583,26 @@ export class OrcaRuntimeService { private async refreshPtyWorktreeRecordsFromController( resolvedWorktrees: ResolvedWorktree[], targetWorktreeId: string | null = null, - deadline?: number + deadline?: number, + signal?: AbortSignal ): Promise | null> { - const inventory = await this.refreshPtyWorktreeRecordsWithControllerInventory( - resolvedWorktrees, - targetWorktreeId, - deadline - ) - return inventory ? new Set(inventory.livePtyIds) : null + this.ptyLivenessRefreshInProgress += 1 + try { + const inventory = await this.refreshPtyWorktreeRecordsWithControllerInventory( + resolvedWorktrees, + targetWorktreeId, + deadline, + undefined, + false, + signal + ) + if (inventory) { + this.ptyLivenessRefreshRequired = false + } + return inventory ? new Set(inventory.livePtyIds) : null + } finally { + this.ptyLivenessRefreshInProgress -= 1 + } } private async refreshPtyWorktreeRecordsWithControllerInventory( @@ -35552,7 +35610,8 @@ export class OrcaRuntimeService { targetWorktreeId: string | null = null, deadline?: number, connectionId?: string | null, - retryStale = false + retryStale = false, + signal?: AbortSignal ): Promise { if (targetWorktreeId === FLOATING_TERMINAL_WORKTREE_ID) { const targetedLiveness = this.refreshFloatingWorkspacePtyLiveness() @@ -35585,7 +35644,8 @@ export class OrcaRuntimeService { // never answers still leaves the aggregate time to return the providers that did // — expiring at the same instant would discard the whole inventory instead. const providerListOpts = { - deadlineMs: Date.now() + Math.max(1, listBudgetMs - PTY_CONTROLLER_LIST_PROVIDER_MARGIN_MS) + deadlineMs: Date.now() + Math.max(1, listBudgetMs - PTY_CONTROLLER_LIST_PROVIDER_MARGIN_MS), + ...(signal ? { signal } : {}) } const processInventory = connectionId === undefined && this.ptyController.listProcessesWithHostScope @@ -35635,7 +35695,8 @@ export class OrcaRuntimeService { targetWorktreeId, deadline, connectionId, - true + true, + signal ) } return null diff --git a/src/main/runtime/rpc/methods/worktree-catalog-methods.ts b/src/main/runtime/rpc/methods/worktree-catalog-methods.ts index 2010a219b7c..d2c76ae1ee8 100644 --- a/src/main/runtime/rpc/methods/worktree-catalog-methods.ts +++ b/src/main/runtime/rpc/methods/worktree-catalog-methods.ts @@ -12,13 +12,15 @@ export const WORKTREE_CATALOG_METHODS: RpcMethod[] = [ name: 'worktree.ps', params: WorktreePsParams, handler: async (params, context) => { - const result = await context.runtime.getWorktreePs( - params.limit, - supportsWorktreeVisibilitySourceDefaults( - context, - params.supportsWorktreeVisibilitySourceDefaults - ) + const supportsSourceDefaults = supportsWorktreeVisibilitySourceDefaults( + context, + params.supportsWorktreeVisibilitySourceDefaults ) + const result = context.signal + ? await context.runtime.getWorktreePs(params.limit, supportsSourceDefaults, { + signal: context.signal + }) + : await context.runtime.getWorktreePs(params.limit, supportsSourceDefaults) // Why: callers that never send the field get the byte-exact legacy response. return params.afterSnapshotId === undefined ? result