mirror of
https://github.com/stablyai/orca.git
synced 2026-10-06 08:02:28 +00:00
perf(wsl): bound guest inventory refresh and correlation
This commit is contained in:
@@ -200,7 +200,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))
|
||||
)
|
||||
|
||||
@@ -1,10 +1,21 @@
|
||||
import type { IPtyProvider } from '../../../providers/types'
|
||||
import { LocalPtyProvider } from '../../../providers/local-pty-provider'
|
||||
import type { PtyProcessInfo } from '../../../providers/pty-process-info'
|
||||
import { parseAppSshPtyId } from '../../../providers/ssh-pty-id'
|
||||
import {
|
||||
LOCAL_EXECUTION_HOST_ID,
|
||||
toSshExecutionHostId,
|
||||
type ExecutionHostId
|
||||
} from '../../../../shared/execution-host'
|
||||
import { ptyOwnership } from '../provider/ownership-state'
|
||||
import { ptySizes } from '../delivery/visibility-state'
|
||||
import { rendererSerializerReadiness } from '../pane/serializer-state'
|
||||
import { getProviderForPty, localProvider } from '../provider/registry'
|
||||
import {
|
||||
getProvider,
|
||||
getProviderForPty,
|
||||
localProvider,
|
||||
registeredPtyProviders
|
||||
} from '../provider/registry'
|
||||
import { inspectPtyProviderProcess } from '../../../providers/pty-process-inspection'
|
||||
import type { PtyRuntimeControllerDeps } from './controller-deps'
|
||||
import { agentSessionPtyWriteGate } from '../../../runtime/agent-session-pty-write-gate'
|
||||
@@ -206,6 +217,68 @@ export function hasPtyFromRuntimeController(
|
||||
}
|
||||
}
|
||||
|
||||
function markSshInventoryUnverifiable(
|
||||
runtime: PtyRuntimeControllerDeps['runtime'],
|
||||
connectionId: string,
|
||||
error: unknown
|
||||
): void {
|
||||
const reason = error instanceof Error ? error.message : String(error)
|
||||
for (const [ptyId, ownerConnectionId] of ptyOwnership) {
|
||||
if (ownerConnectionId === connectionId) {
|
||||
runtime?.markPtyLivenessUnverifiable?.(ptyId, reason)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
export async function listProcessesWithHostScopeFromRuntimeController(
|
||||
deps: PtyRuntimeControllerDeps,
|
||||
opts?: { deadlineMs?: number; signal?: AbortSignal }
|
||||
): Promise<{ processes: PtyProcessInfo[]; hostIds: ExecutionHostId[] }> {
|
||||
const providerSessions = await Promise.all(
|
||||
registeredPtyProviders().map(async ({ provider, connectionId }) => {
|
||||
const hostId: ExecutionHostId = connectionId
|
||||
? toSshExecutionHostId(connectionId)
|
||||
: LOCAL_EXECUTION_HOST_ID
|
||||
try {
|
||||
return {
|
||||
processes: await provider.listProcesses(opts),
|
||||
hostId
|
||||
}
|
||||
} catch (error) {
|
||||
if (!connectionId) {
|
||||
throw error
|
||||
}
|
||||
markSshInventoryUnverifiable(deps.runtime, connectionId, error)
|
||||
return null
|
||||
}
|
||||
})
|
||||
)
|
||||
const respondingSessions = providerSessions.filter((session) => session !== null)
|
||||
return {
|
||||
processes: respondingSessions.flatMap((session) => session.processes),
|
||||
hostIds: respondingSessions.map((session) => session.hostId)
|
||||
}
|
||||
}
|
||||
|
||||
export async function listProcessesFromRuntimeController(
|
||||
deps: PtyRuntimeControllerDeps,
|
||||
connectionId?: string | null,
|
||||
opts?: { deadlineMs?: number; signal?: AbortSignal }
|
||||
) {
|
||||
if (connectionId === null) {
|
||||
return localProvider.listProcesses(opts)
|
||||
}
|
||||
if (connectionId !== undefined) {
|
||||
try {
|
||||
return await getProvider(connectionId).listProcesses(opts)
|
||||
} catch (error) {
|
||||
markSshInventoryUnverifiable(deps.runtime, connectionId, error)
|
||||
throw error
|
||||
}
|
||||
}
|
||||
return (await listProcessesWithHostScopeFromRuntimeController(deps, opts)).processes
|
||||
}
|
||||
|
||||
export function resizePtyFromRuntimeController(ptyId: string, cols: number, rows: number): boolean {
|
||||
try {
|
||||
getProviderForPty(ptyId).resize(ptyId, cols, rows)
|
||||
@@ -240,7 +313,7 @@ export function getSizeFromRuntimeController(ptyId: string) {
|
||||
|
||||
export async function serializeProviderBufferFromRuntimeController(
|
||||
ptyId: string,
|
||||
opts?: { scrollbackRows?: number }
|
||||
opts?: { scrollbackRows?: number; altScreenForcesZeroRows?: boolean }
|
||||
) {
|
||||
try {
|
||||
// Why: restored daemon PTYs can be live while their desktop pane is unmounted; query the provider model so phone-local navigation works.
|
||||
|
||||
@@ -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'
|
||||
|
||||
@@ -232,6 +232,7 @@ export type IPtyProvider = {
|
||||
// Why: deadlineMs bounds the underlying RPC exactly like shutdown's deadlineMs.
|
||||
listProcesses(opts?: {
|
||||
deadlineMs?: number
|
||||
signal?: AbortSignal
|
||||
includeForegroundProcessEvidence?: boolean
|
||||
}): Promise<PtyProcessInfo[]>
|
||||
getDefaultShell(): Promise<string>
|
||||
|
||||
@@ -26,8 +26,17 @@ import { spawnWithTerminalRuntimeRepair, type TerminalRepairHook } from './ssh-p
|
||||
import { createSshPtyProviderRpcOperations } from './ssh-pty-provider-rpc-operations'
|
||||
|
||||
// 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. */
|
||||
@@ -268,13 +277,14 @@ export class SshPtyProvider implements IPtyProvider {
|
||||
async listProcesses(opts?: {
|
||||
deadlineMs?: number
|
||||
includeForegroundProcessEvidence?: boolean
|
||||
signal?: AbortSignal
|
||||
}): Promise<PtyProcessInfo[]> {
|
||||
const result = await this.mux.request(
|
||||
'pty.listProcesses',
|
||||
opts?.includeForegroundProcessEvidence === undefined
|
||||
? undefined
|
||||
: { includeForegroundProcessEvidence: opts.includeForegroundProcessEvidence },
|
||||
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
|
||||
|
||||
+44563
-55
File diff suppressed because it is too large
Load Diff
@@ -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