mirror of
https://github.com/stablyai/orca.git
synced 2026-09-30 00:03:15 +00:00
perf(wsl): bound guest inventory refresh and correlation
This commit is contained in:
@@ -197,7 +197,10 @@ export class DaemonPtyRouter implements IPtyProvider {
|
||||
await this.current.revive(state)
|
||||
}
|
||||
|
||||
async listProcesses(opts?: { deadlineMs?: number }): Promise<PtyProcessInfo[]> {
|
||||
async listProcesses(opts?: {
|
||||
deadlineMs?: number
|
||||
signal?: AbortSignal
|
||||
}): Promise<PtyProcessInfo[]> {
|
||||
// Why: runtime exact-stop/liveness flows must fail closed if any adapter
|
||||
// cannot provide a trustworthy process list.
|
||||
const results = await Promise.all(
|
||||
|
||||
@@ -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<PtyProcessInfo[]> {
|
||||
async listProcesses(opts?: {
|
||||
deadlineMs?: number
|
||||
signal?: AbortSignal
|
||||
}): Promise<PtyProcessInfo[]> {
|
||||
// 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<void> => {
|
||||
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<string, ReturnType<typeof createWslGuestProcessIndexes>>()
|
||||
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'
|
||||
|
||||
@@ -206,7 +206,10 @@ export class DegradedDaemonPtyProvider implements IPtyProvider {
|
||||
await this.fallback.revive(state)
|
||||
}
|
||||
|
||||
async listProcesses(opts?: { deadlineMs?: number }): Promise<PtyProcessInfo[]> {
|
||||
async listProcesses(opts?: {
|
||||
deadlineMs?: number
|
||||
signal?: AbortSignal
|
||||
}): Promise<PtyProcessInfo[]> {
|
||||
const results = await Promise.all(
|
||||
this.allProviders().map((provider) => provider.listProcesses(opts))
|
||||
)
|
||||
|
||||
@@ -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 {
|
||||
|
||||
@@ -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<void>
|
||||
@@ -132,6 +133,7 @@ export function clearPtyState(id: string): void {
|
||||
ptyTerminationMode.delete(id)
|
||||
ptyReportsChildExitStatus.delete(id)
|
||||
ptyPhysicalExits.delete(id)
|
||||
resetWslGuestProcessInventory()
|
||||
}
|
||||
|
||||
/**
|
||||
|
||||
@@ -136,8 +136,8 @@ export class LocalPtyProvider implements IPtyProvider {
|
||||
/* re-spawning handles local revival */
|
||||
}
|
||||
|
||||
listProcesses(): Promise<PtyProcessInfo[]> {
|
||||
return listLocalPtyProcesses()
|
||||
listProcesses(opts?: { deadlineMs?: number; signal?: AbortSignal }): Promise<PtyProcessInfo[]> {
|
||||
return listLocalPtyProcesses(opts)
|
||||
}
|
||||
|
||||
getDefaultShell(): Promise<string> {
|
||||
|
||||
@@ -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<PtyProcessInfo[]> {
|
||||
export async function listLocalPtyProcesses(opts?: {
|
||||
deadlineMs?: number
|
||||
signal?: AbortSignal
|
||||
}): Promise<PtyProcessInfo[]> {
|
||||
const entries = Array.from(ptyProcesses.entries())
|
||||
const evidenceEpoch = Date.now()
|
||||
const wslByDistro = new Map<string, string[]>()
|
||||
@@ -137,11 +142,31 @@ export async function listLocalPtyProcesses(): Promise<PtyProcessInfo[]> {
|
||||
}
|
||||
}
|
||||
const inventories = new Map<string, WslGuestProcessInventoryRead>()
|
||||
const distros = [...wslByDistro.keys()]
|
||||
let nextDistroIndex = 0
|
||||
const readNextDistro = async (): Promise<void> => {
|
||||
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<string, ReturnType<typeof createWslGuestProcessIndexes>>()
|
||||
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<PtyProcessInfo[]> {
|
||||
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'
|
||||
|
||||
@@ -222,7 +222,7 @@ export type IPtyProvider = {
|
||||
serialize(ids: string[]): Promise<string>
|
||||
revive(state: string): Promise<void>
|
||||
// Why: deadlineMs bounds the underlying RPC exactly like shutdown's deadlineMs.
|
||||
listProcesses(opts?: { deadlineMs?: number }): Promise<PtyProcessInfo[]>
|
||||
listProcesses(opts?: { deadlineMs?: number; signal?: AbortSignal }): Promise<PtyProcessInfo[]>
|
||||
getDefaultShell(): Promise<string>
|
||||
getProfiles(): Promise<{ name: string; path: string }[]>
|
||||
onData(callback: (payload: PtyDataEvent) => void): () => void
|
||||
|
||||
@@ -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<PtyProcessInfo[]> {
|
||||
async listProcesses(opts?: {
|
||||
deadlineMs?: number
|
||||
signal?: AbortSignal
|
||||
}): Promise<PtyProcessInfo[]> {
|
||||
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) {
|
||||
|
||||
@@ -11,14 +11,54 @@ export type WslGuestForegroundResolution =
|
||||
| { status: 'live'; processName: string | null; anchor: WslGuestProcessAnchor }
|
||||
| { status: 'unverifiable'; reason: string }
|
||||
|
||||
export type WslGuestProcessIndexes = {
|
||||
byPid: ReadonlyMap<number, WslGuestProcessRow>
|
||||
byForegroundGroup: ReadonlyMap<string, readonly WslGuestProcessRow[]>
|
||||
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<number, WslGuestProcessRow>()
|
||||
const groups = new Map<string, WslGuestProcessRow[]>()
|
||||
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<number>()
|
||||
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' }
|
||||
|
||||
@@ -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({
|
||||
|
||||
@@ -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<WslGuestProcessInventoryRead>
|
||||
run?: (
|
||||
distro: string,
|
||||
opts?: { deadlineMs?: number; signal?: AbortSignal }
|
||||
) => Promise<WslGuestProcessInventoryRead>
|
||||
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<WslGuestProcessInventoryRead>
|
||||
read: (
|
||||
distro: string,
|
||||
opts?: WslGuestProcessInventoryReadOptions
|
||||
) => Promise<WslGuestProcessInventoryRead>
|
||||
reset: () => void
|
||||
} {
|
||||
const now = deps.now ?? (() => Date.now())
|
||||
const ttlMs = deps.ttlMs ?? INVENTORY_TTL_MS
|
||||
const cached = new Map<string, { value: WslGuestProcessInventoryRead; at: number }>()
|
||||
const inFlight = new Map<string, Promise<WslGuestProcessInventoryRead>>()
|
||||
let resetGeneration = 0
|
||||
const run = deps.run ?? runWslGuestProcessInventory
|
||||
|
||||
const read = (distro: string): Promise<WslGuestProcessInventoryRead> => {
|
||||
const read = (
|
||||
distro: string,
|
||||
opts?: WslGuestProcessInventoryReadOptions
|
||||
): Promise<WslGuestProcessInventoryRead> => {
|
||||
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<WslGuestProcessInventoryRead> {
|
||||
async function runWslGuestProcessInventory(
|
||||
distro: string,
|
||||
opts?: WslGuestProcessInventoryReadOptions
|
||||
): Promise<WslGuestProcessInventoryRead> {
|
||||
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<WslGuestProc
|
||||
WSLENV: wslenvEntries.join(':'),
|
||||
[CAPTURE_NONCE_ENV]: captureNonce
|
||||
},
|
||||
timeoutMs: INVENTORY_TIMEOUT_MS,
|
||||
timeoutMs:
|
||||
opts?.deadlineMs === undefined
|
||||
? INVENTORY_TIMEOUT_MS
|
||||
: Math.max(1, Math.min(INVENTORY_TIMEOUT_MS, opts.deadlineMs - Date.now())),
|
||||
...(opts?.signal ? { signal: opts.signal } : {}),
|
||||
maxOutputBytes: INVENTORY_MAX_OUTPUT_BYTES
|
||||
})
|
||||
} catch {
|
||||
return { status: 'unverifiable', reason: 'wsl_unavailable' }
|
||||
}
|
||||
if (result.timedOut) {
|
||||
if (result.timedOut || opts?.signal?.aborted) {
|
||||
return { status: 'unverifiable', reason: 'capture_timed_out' }
|
||||
}
|
||||
if (result.code === 127) {
|
||||
@@ -197,11 +256,14 @@ async function runWslGuestProcessInventory(distro: string): Promise<WslGuestProc
|
||||
const defaultReader = createWslGuestProcessInventoryReader()
|
||||
|
||||
export function readWslGuestProcessInventory(
|
||||
distro: string
|
||||
distro: string,
|
||||
opts?: WslGuestProcessInventoryReadOptions
|
||||
): Promise<WslGuestProcessInventoryRead> {
|
||||
return defaultReader.read(distro)
|
||||
return defaultReader.read(distro, opts)
|
||||
}
|
||||
|
||||
export function resetWslGuestProcessInventoryForTests(): void {
|
||||
export function resetWslGuestProcessInventory(): void {
|
||||
defaultReader.reset()
|
||||
}
|
||||
|
||||
export const resetWslGuestProcessInventoryForTests = resetWslGuestProcessInventory
|
||||
|
||||
@@ -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 = {
|
||||
|
||||
@@ -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<PtyProcessInfo[]>
|
||||
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<string, TrackedPtyLivenessVerdict>()
|
||||
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<string>()
|
||||
private wslDistroByPtyId = new Map<string, string>()
|
||||
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>) => 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<string, RuntimeWorktreePsSummary>()
|
||||
@@ -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<Set<string> | 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<PtyControllerInventory | null> {
|
||||
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
|
||||
|
||||
@@ -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
|
||||
|
||||
Reference in New Issue
Block a user