mirror of
https://github.com/stablyai/orca.git
synced 2026-09-29 16:02:50 +00:00
fix(computer-use): recover the runtime host instead of latching it off
Review follow-up on the persistent Windows computer-use helper. A helper that died before producing a line set an unavailable flag nothing ever cleared, and the client then dropped the host for the life of the session. One transient bad spawn — a Defender scan, a locked CSC temp directory — silently restored the per-click powershell.exe burst and per-operation MSIL emission this work exists to remove, with computer use still working so nothing looked wrong. Start failures are now retried, then cool down for 60s, then re-probed; the client keeps the host so it can come back. Repeated post-answer crashes cool down too, and a single reply no longer clears the failure count. The one-shot bridge decided its execution-policy retry from a message that fell back to stdout, so a window title containing "SecurityError" could replay a non-idempotent operation — a double click, keystroke or paste — and stick the session on Bypass. The retry now requires empty stdout and a matching stderr. Serve-mode replies carry an echoed request id. Without one a single stray stdout line would make every later response answer the previous request, acting on stale element indexes with no error raised; a mismatch now kills the child. Non-JSON noise is ignored rather than counted as the helper having answered. Also: warnings reach the main process over the sidecar's IPC channel rather than its piped, unread stdio; the child is watched on close rather than exit; dispose latches so a queued request cannot respawn during teardown; and the host is split into a serve channel and an availability policy to stay under max-lines.
This commit is contained in:
@@ -7,6 +7,9 @@ param(
|
||||
)
|
||||
|
||||
$ErrorActionPreference = "Stop"
|
||||
# Progress records render to the host, which in serve mode is a pipe carrying
|
||||
# one JSON response per line; a stray record would desynchronise the stream.
|
||||
$ProgressPreference = "SilentlyContinue"
|
||||
$utf8NoBom = New-Object System.Text.UTF8Encoding $false
|
||||
[Console]::InputEncoding = $utf8NoBom
|
||||
[Console]::OutputEncoding = $utf8NoBom
|
||||
@@ -1324,12 +1327,20 @@ function Invoke-OrcaServeLoop {
|
||||
$line = [Console]::In.ReadLine()
|
||||
if ($null -eq $line) { break }
|
||||
if ([string]::IsNullOrWhiteSpace($line)) { continue }
|
||||
$requestId = $null
|
||||
try {
|
||||
$json = ConvertTo-Json (Invoke-OrcaOperation ($line | ConvertFrom-Json)) -Depth 100 -Compress
|
||||
$operation = $line | ConvertFrom-Json
|
||||
$requestId = $operation.requestId
|
||||
$response = Invoke-OrcaOperation $operation
|
||||
} catch {
|
||||
$json = ConvertTo-Json ([pscustomobject]@{ ok = $false; error = [string]$_.Exception.Message }) -Depth 100 -Compress
|
||||
$response = [pscustomobject]@{ ok = $false; error = [string]$_.Exception.Message }
|
||||
}
|
||||
[Console]::Out.WriteLine($json)
|
||||
# Echoed so the caller can prove which request a line answers; a reply it
|
||||
# cannot match is a desynchronised stream, not a usable response.
|
||||
if ($null -ne $requestId) {
|
||||
$response | Add-Member -NotePropertyName requestId -NotePropertyValue $requestId -Force
|
||||
}
|
||||
[Console]::Out.WriteLine((ConvertTo-Json $response -Depth 100 -Compress))
|
||||
[Console]::Out.Flush()
|
||||
}
|
||||
}
|
||||
|
||||
@@ -0,0 +1,45 @@
|
||||
import { afterEach, describe, expect, it, vi } from 'vitest'
|
||||
import {
|
||||
isComputerSidecarDiagnostic,
|
||||
reportComputerDiagnostic
|
||||
} from './computer-sidecar-diagnostics'
|
||||
|
||||
describe('computer sidecar diagnostics', () => {
|
||||
const originalSend = process.send
|
||||
|
||||
afterEach(() => {
|
||||
process.send = originalSend
|
||||
vi.restoreAllMocks()
|
||||
})
|
||||
|
||||
it('sends over IPC when running inside the sidecar', () => {
|
||||
const send = vi.fn((_message: unknown) => true)
|
||||
process.send = send as unknown as typeof process.send
|
||||
const console_ = vi.spyOn(console, 'warn').mockImplementation(() => {})
|
||||
|
||||
reportComputerDiagnostic('fell back to Bypass')
|
||||
|
||||
// The sidecar's stdout is piped and never read, so this must not go there.
|
||||
expect(console_).not.toHaveBeenCalled()
|
||||
expect(send).toHaveBeenCalledWith({
|
||||
kind: 'computer-sidecar-diagnostic',
|
||||
message: 'fell back to Bypass'
|
||||
})
|
||||
expect(isComputerSidecarDiagnostic(send.mock.calls[0][0])).toBe(true)
|
||||
})
|
||||
|
||||
it('logs directly when there is no IPC channel', () => {
|
||||
process.send = undefined
|
||||
const console_ = vi.spyOn(console, 'warn').mockImplementation(() => {})
|
||||
|
||||
reportComputerDiagnostic('fell back to Bypass')
|
||||
|
||||
expect(console_).toHaveBeenCalledWith('[computer-use] fell back to Bypass')
|
||||
})
|
||||
|
||||
it('does not mistake a sidecar response for a diagnostic', () => {
|
||||
expect(isComputerSidecarDiagnostic({ id: 1, ok: true, result: {} })).toBe(false)
|
||||
expect(isComputerSidecarDiagnostic({ kind: 'computer-sidecar-diagnostic' })).toBe(false)
|
||||
expect(isComputerSidecarDiagnostic(null)).toBe(false)
|
||||
})
|
||||
})
|
||||
@@ -0,0 +1,38 @@
|
||||
/**
|
||||
* Warnings from the computer-use provider, routed to somewhere a human sees.
|
||||
*
|
||||
* Why not `console.warn`: the provider runs inside the forked sidecar, which
|
||||
* `sidecar-client.ts` starts with piped stdio that nothing ever reads. Anything
|
||||
* written there is discarded — including the only signal that a machine has
|
||||
* fallen back to `-ExecutionPolicy Bypass`, a state that persists for the
|
||||
* session. The sidecar has an IPC channel already, so the warning takes it.
|
||||
*/
|
||||
export type ComputerSidecarDiagnostic = {
|
||||
kind: 'computer-sidecar-diagnostic'
|
||||
message: string
|
||||
}
|
||||
|
||||
const DIAGNOSTIC_KIND = 'computer-sidecar-diagnostic'
|
||||
|
||||
export function isComputerSidecarDiagnostic(
|
||||
message: unknown
|
||||
): message is ComputerSidecarDiagnostic {
|
||||
if (!message || typeof message !== 'object') {
|
||||
return false
|
||||
}
|
||||
const record = message as Record<string, unknown>
|
||||
return record.kind === DIAGNOSTIC_KIND && typeof record.message === 'string'
|
||||
}
|
||||
|
||||
export function reportComputerDiagnostic(message: string): void {
|
||||
if (process.send) {
|
||||
process.send({ kind: DIAGNOSTIC_KIND, message } satisfies ComputerSidecarDiagnostic)
|
||||
return
|
||||
}
|
||||
logComputerDiagnostic(message)
|
||||
}
|
||||
|
||||
/** The main-process end: how a sidecar's forwarded diagnostic is printed. */
|
||||
export function logComputerDiagnostic(message: string): void {
|
||||
console.warn(`[computer-use] ${message}`)
|
||||
}
|
||||
@@ -1,5 +1,6 @@
|
||||
import { execFile } from 'node:child_process'
|
||||
import { windowsPowerShellPath } from '../../shared/child-process/windows-system-binary'
|
||||
import { reportComputerDiagnostic } from './computer-sidecar-diagnostics'
|
||||
import { RuntimeClientError } from './runtime-client-error'
|
||||
import type { DesktopScriptPlatform } from './desktop-script-provider-paths'
|
||||
import {
|
||||
@@ -18,7 +19,7 @@ export async function execBridge(
|
||||
operationPath: string
|
||||
): Promise<{ stdout: string; stderr: string }> {
|
||||
if (platform !== 'windows') {
|
||||
return await runBridgeProcess('python3', [scriptPath, operationPath])
|
||||
return await mapped(runBridgeProcess('python3', [scriptPath, operationPath]))
|
||||
}
|
||||
const command = windowsPowerShellPath()
|
||||
try {
|
||||
@@ -27,20 +28,61 @@ export async function execBridge(
|
||||
windowsPowerShellRuntimeArgs(scriptPath, PREFERRED_WINDOWS_EXECUTION_POLICY, [operationPath])
|
||||
)
|
||||
} catch (error) {
|
||||
// The script never loaded under a blocking policy, so re-running is safe.
|
||||
if (!isExecutionPolicyBlocked(error instanceof Error ? error.message : String(error))) {
|
||||
throw error
|
||||
if (!isPolicyBlockedStart(error)) {
|
||||
throw error instanceof BridgeProcessFailure ? error.mapped : error
|
||||
}
|
||||
console.warn(
|
||||
`[computer-use] bridge start blocked at ${PREFERRED_WINDOWS_EXECUTION_POLICY}; retrying once with ${FALLBACK_WINDOWS_EXECUTION_POLICY}`
|
||||
reportComputerDiagnostic(
|
||||
`bridge start blocked at ${PREFERRED_WINDOWS_EXECUTION_POLICY}; retrying once with ${FALLBACK_WINDOWS_EXECUTION_POLICY}`
|
||||
)
|
||||
return await runBridgeProcess(
|
||||
command,
|
||||
windowsPowerShellRuntimeArgs(scriptPath, FALLBACK_WINDOWS_EXECUTION_POLICY, [operationPath])
|
||||
return await mapped(
|
||||
runBridgeProcess(
|
||||
command,
|
||||
windowsPowerShellRuntimeArgs(scriptPath, FALLBACK_WINDOWS_EXECUTION_POLICY, [operationPath])
|
||||
)
|
||||
)
|
||||
}
|
||||
}
|
||||
|
||||
/** Unwrap the raw-stream carrier back into the error callers expect. */
|
||||
async function mapped(
|
||||
run: Promise<{ stdout: string; stderr: string }>
|
||||
): Promise<{ stdout: string; stderr: string }> {
|
||||
try {
|
||||
return await run
|
||||
} catch (error) {
|
||||
throw error instanceof BridgeProcessFailure ? error.mapped : error
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* Only a run that produced no stdout at all may be replayed.
|
||||
*
|
||||
* Why the stdout guard: operations are not idempotent, and the response embeds
|
||||
* window titles and element names. Matching the policy pattern against a
|
||||
* snapshot that merely contains the word "SecurityError" would replay the
|
||||
* operation — a second click, a second keystroke, a second paste. A helper
|
||||
* blocked by the execution policy never reaches its first line of output.
|
||||
*/
|
||||
function isPolicyBlockedStart(error: unknown): error is BridgeProcessFailure {
|
||||
return (
|
||||
error instanceof BridgeProcessFailure &&
|
||||
!error.stdout.trim() &&
|
||||
isExecutionPolicyBlocked(error.stderr)
|
||||
)
|
||||
}
|
||||
|
||||
/** Carries the raw streams so the retry decision does not read a mapped message. */
|
||||
class BridgeProcessFailure extends Error {
|
||||
constructor(
|
||||
readonly stdout: string,
|
||||
readonly stderr: string,
|
||||
readonly mapped: RuntimeClientError
|
||||
) {
|
||||
super(mapped.message)
|
||||
this.name = 'BridgeProcessFailure'
|
||||
}
|
||||
}
|
||||
|
||||
function runBridgeProcess(
|
||||
command: string,
|
||||
args: readonly string[]
|
||||
@@ -108,11 +150,10 @@ function runBridgeProcess(
|
||||
(error, stdout, stderr) => {
|
||||
if (error) {
|
||||
const message = stderr.trim() || stdout.trim() || error.message
|
||||
finish(
|
||||
error.killed
|
||||
? new RuntimeClientError('action_timeout', message)
|
||||
: mapBridgeError(message)
|
||||
)
|
||||
const mapped = error.killed
|
||||
? new RuntimeClientError('action_timeout', message)
|
||||
: mapBridgeError(message)
|
||||
finish(new BridgeProcessFailure(stdout, stderr, mapped))
|
||||
return
|
||||
}
|
||||
finish(null, { stdout, stderr })
|
||||
|
||||
@@ -53,7 +53,10 @@ export class DesktopScriptProviderClient {
|
||||
constructor(
|
||||
private readonly platform: DesktopScriptPlatform = requiredPlatform(),
|
||||
private readonly scriptPath: string = requiredScriptPath(),
|
||||
private runtimeHost: DesktopScriptRuntimeHost | null = defaultRuntimeHost(platform, scriptPath)
|
||||
private readonly runtimeHost: DesktopScriptRuntimeHost | null = defaultRuntimeHost(
|
||||
platform,
|
||||
scriptPath
|
||||
)
|
||||
) {}
|
||||
|
||||
shutdown(): void {
|
||||
@@ -212,11 +215,11 @@ export class DesktopScriptProviderClient {
|
||||
return checkedBridgeResponse(await host.request(request), '')
|
||||
} catch (error) {
|
||||
// Only a helper that cannot start falls back; operation errors surface.
|
||||
// The host is kept: it re-probes after its cooldown, so a transient bad
|
||||
// spawn cannot strand the session on one powershell.exe per operation.
|
||||
if (!isRuntimeHostUnavailable(error)) {
|
||||
throw error
|
||||
}
|
||||
this.runtimeHost = null
|
||||
host.dispose()
|
||||
}
|
||||
}
|
||||
return await this.callOneShotBridge(request)
|
||||
|
||||
@@ -47,7 +47,7 @@ describe('desktop script provider runtime host routing', () => {
|
||||
expectDesktopProviderSubprocessStartCount(0)
|
||||
})
|
||||
|
||||
it('degrades to the one-shot bridge when the runtime host cannot start', async () => {
|
||||
it('degrades to the one-shot bridge for the operations a host cannot serve', async () => {
|
||||
const request = vi.fn(async () => {
|
||||
throw new RuntimeClientError('runtime_host_unavailable', 'could not start')
|
||||
})
|
||||
@@ -58,14 +58,36 @@ describe('desktop script provider runtime host routing', () => {
|
||||
const client = await createDesktopScriptProviderClient('windows', 'C:\\runtime.ps1', host)
|
||||
|
||||
await expect(client.listApps()).resolves.toMatchObject({ apps: [{ pid: 42 }] })
|
||||
expect(dispose).toHaveBeenCalled()
|
||||
|
||||
// The host is dropped for the session rather than re-probed per operation.
|
||||
await client.listApps()
|
||||
expect(request).toHaveBeenCalledTimes(1)
|
||||
|
||||
// The host is kept and asked again: it owns its own cooldown, so one bad
|
||||
// spawn must not stand the session down to a powershell.exe per click.
|
||||
expect(request).toHaveBeenCalledTimes(2)
|
||||
expect(dispose).not.toHaveBeenCalled()
|
||||
expectDesktopProviderSubprocessStartCount(2)
|
||||
})
|
||||
|
||||
it('returns to the runtime host once it recovers', async () => {
|
||||
let healthy = false
|
||||
const request = vi.fn(async () => {
|
||||
if (!healthy) {
|
||||
throw new RuntimeClientError('runtime_host_unavailable', 'could not start')
|
||||
}
|
||||
return { ok: true, apps: [] } as BridgeResponse
|
||||
})
|
||||
const { host } = fakeRuntimeHost(request)
|
||||
mockBridgeResponse({ ok: true, apps: [] })
|
||||
|
||||
const client = await createDesktopScriptProviderClient('windows', 'C:\\runtime.ps1', host)
|
||||
|
||||
await client.listApps()
|
||||
expectDesktopProviderSubprocessStartCount(1)
|
||||
|
||||
healthy = true
|
||||
await expect(client.listApps()).resolves.toEqual({ apps: [] })
|
||||
expectDesktopProviderSubprocessStartCount(1)
|
||||
})
|
||||
|
||||
it('runs the one-shot bridge under RemoteSigned and falls back to Bypass once', async () => {
|
||||
mockBridgeProcessFailure(POLICY_STDERR)
|
||||
mockBridgeResponse({ ok: true, apps: [] })
|
||||
@@ -74,6 +96,7 @@ describe('desktop script provider runtime host routing', () => {
|
||||
|
||||
await expect(client.listApps()).resolves.toEqual({ apps: [] })
|
||||
expectDesktopProviderSubprocessStartCount(2)
|
||||
expect(bridgeProcessArgs(0)).toContain('-NoLogo')
|
||||
expect(bridgeProcessArgs(0)).toContain('RemoteSigned')
|
||||
expect(bridgeProcessArgs(0)).not.toContain('Bypass')
|
||||
expect(bridgeProcessArgs(1)).toContain('Bypass')
|
||||
@@ -88,6 +111,23 @@ describe('desktop script provider runtime host routing', () => {
|
||||
expectDesktopProviderSubprocessStartCount(1)
|
||||
})
|
||||
|
||||
it('never replays an operation whose own output merely mentions a policy error', async () => {
|
||||
// Window titles and element names are user-controlled text that lands in
|
||||
// stdout; matching them would double a click, a keystroke or a paste.
|
||||
mockBridgeProcessFailure({
|
||||
stdout: JSON.stringify({
|
||||
ok: true,
|
||||
snapshot: { windowTitle: 'SecurityError - UnauthorizedAccess.log - Notepad' }
|
||||
}),
|
||||
stderr: ''
|
||||
})
|
||||
|
||||
const client = await createDesktopScriptProviderClient('windows', 'C:\\runtime.ps1')
|
||||
|
||||
await expect(client.listApps()).rejects.toBeInstanceOf(Error)
|
||||
expectDesktopProviderSubprocessStartCount(1)
|
||||
})
|
||||
|
||||
it('keeps Linux on the one-shot python bridge with no execution policy flags', async () => {
|
||||
mockBridgeResponse({ ok: true, apps: [] })
|
||||
|
||||
|
||||
@@ -80,10 +80,11 @@ export function mockBridgeResponse(
|
||||
})
|
||||
}
|
||||
|
||||
export function mockBridgeProcessFailure(stderr: string): void {
|
||||
export function mockBridgeProcessFailure(streams: string | { stdout?: string; stderr?: string }) {
|
||||
const { stdout = '', stderr = '' } = typeof streams === 'string' ? { stderr: streams } : streams
|
||||
execFileMock.mockImplementationOnce((_command, _args, _options, callback) => {
|
||||
const done = callback as (error: Error | null, stdout: string, stderr: string) => void
|
||||
done(new Error('Command failed'), '', stderr)
|
||||
done(new Error('Command failed'), stdout, stderr)
|
||||
return null as never
|
||||
})
|
||||
}
|
||||
|
||||
@@ -101,6 +101,8 @@ export type BridgeWindow = {
|
||||
|
||||
export type BridgeResponse = {
|
||||
ok: boolean
|
||||
/** Echo of BridgeRequest.requestId; set only on the persistent serve path. */
|
||||
requestId?: number
|
||||
error?: string
|
||||
capabilities?: ComputerProviderCapabilities
|
||||
apps?: {
|
||||
@@ -122,6 +124,8 @@ export type BridgeResponse = {
|
||||
|
||||
export type BridgeRequest = {
|
||||
tool: string
|
||||
/** Correlates a serve-mode reply with its request; the one-shot path omits it. */
|
||||
requestId?: number
|
||||
app?: string
|
||||
element?: BridgeElement
|
||||
fromElement?: BridgeElement
|
||||
|
||||
@@ -0,0 +1,109 @@
|
||||
import {
|
||||
FALLBACK_WINDOWS_EXECUTION_POLICY,
|
||||
PREFERRED_WINDOWS_EXECUTION_POLICY,
|
||||
type WindowsExecutionPolicy
|
||||
} from './windows-powershell-execution-policy'
|
||||
|
||||
/**
|
||||
* Consecutive child failures before the helper is believed dead, and how long
|
||||
* the one-shot bridge covers for it afterwards.
|
||||
*
|
||||
* Why not a latch: every plausible cause is transient — a Defender scan touching
|
||||
* the script mid-launch, a locked CSC temp directory failing one `Add-Type`,
|
||||
* momentary memory pressure. Giving up permanently silently restores the
|
||||
* per-click process burst the host exists to remove, and computer use keeps
|
||||
* working throughout, so nothing looks wrong while the MDE signature returns.
|
||||
*/
|
||||
export const MAX_START_ATTEMPTS = 3
|
||||
export const START_FAILURE_COOLDOWN_MS = 60_000
|
||||
|
||||
/**
|
||||
* Whether the persistent helper is currently believed usable, and the execution
|
||||
* policy it should be started under.
|
||||
*
|
||||
* Split from the host so the recovery rules are readable on their own: they are
|
||||
* what stands between a transient bad spawn and a session that silently spends
|
||||
* the rest of its life on one powershell.exe per click.
|
||||
*/
|
||||
export class RuntimeHostAvailability {
|
||||
private policy: WindowsExecutionPolicy = PREFERRED_WINDOWS_EXECUTION_POLICY
|
||||
private retryUnderFallbackPolicy = false
|
||||
private consecutiveFailures = 0
|
||||
private consecutiveSuccesses = 0
|
||||
private cooldownUntil = 0
|
||||
|
||||
constructor(
|
||||
private readonly cooldownMs: number,
|
||||
private readonly now: () => number,
|
||||
/** Public so the host can report its own start attempts to the same sink. */
|
||||
readonly warn: (message: string) => void
|
||||
) {}
|
||||
|
||||
get executionPolicy(): WindowsExecutionPolicy {
|
||||
return this.policy
|
||||
}
|
||||
|
||||
get policyRetryPending(): boolean {
|
||||
return this.retryUnderFallbackPolicy
|
||||
}
|
||||
|
||||
get atPreferredPolicy(): boolean {
|
||||
return this.policy === PREFERRED_WINDOWS_EXECUTION_POLICY
|
||||
}
|
||||
|
||||
/** Milliseconds left before the host may try a helper again; 0 when it may. */
|
||||
remainingCooldown(): number {
|
||||
return Math.max(0, this.cooldownUntil - this.now())
|
||||
}
|
||||
|
||||
requestPolicyRetry(): void {
|
||||
this.retryUnderFallbackPolicy = true
|
||||
}
|
||||
|
||||
escalateExecutionPolicy(): void {
|
||||
this.retryUnderFallbackPolicy = false
|
||||
this.policy = FALLBACK_WINDOWS_EXECUTION_POLICY
|
||||
// Sticky for the session: a genuinely Restricted machine would otherwise pay
|
||||
// a guaranteed failed spawn per operation. Only a helper that produced no
|
||||
// output at all can reach here, so a snapshot cannot talk the host into it.
|
||||
this.warn(
|
||||
`runtime host start blocked at ${PREFERRED_WINDOWS_EXECUTION_POLICY}; using ${FALLBACK_WINDOWS_EXECUTION_POLICY} for the rest of this session`
|
||||
)
|
||||
}
|
||||
|
||||
recordFailure(): void {
|
||||
this.consecutiveSuccesses = 0
|
||||
this.consecutiveFailures++
|
||||
}
|
||||
|
||||
/** True once a helper has died often enough that respawning is just thrash. */
|
||||
get exhausted(): boolean {
|
||||
return this.consecutiveFailures >= MAX_START_ATTEMPTS
|
||||
}
|
||||
|
||||
recordSuccess(): void {
|
||||
this.consecutiveSuccesses++
|
||||
// Why a clean run and not a single reply: a helper that answers one
|
||||
// operation and dies on the next would otherwise reset the count forever,
|
||||
// and respawn once per operation — the exact burst the host removes.
|
||||
if (this.consecutiveSuccesses >= MAX_START_ATTEMPTS) {
|
||||
this.consecutiveFailures = 0
|
||||
}
|
||||
if (this.cooldownUntil === 0) {
|
||||
return
|
||||
}
|
||||
this.cooldownUntil = 0
|
||||
this.warn('runtime host recovered; operations are served by the persistent helper again')
|
||||
}
|
||||
|
||||
enterCooldown(): void {
|
||||
this.cooldownUntil = this.now() + this.cooldownMs
|
||||
this.warn(
|
||||
`runtime host unavailable after ${this.consecutiveFailures} consecutive failures; falling back to one powershell.exe per operation for ${this.cooldownMs}ms`
|
||||
)
|
||||
}
|
||||
|
||||
clearCooldown(): void {
|
||||
this.cooldownUntil = 0
|
||||
}
|
||||
}
|
||||
@@ -1,11 +1,11 @@
|
||||
import { EventEmitter } from 'node:events'
|
||||
import { afterEach, describe, expect, it, vi } from 'vitest'
|
||||
import type { ProcessSpec } from '../../shared/child-process/process-spec'
|
||||
import type { spawnProcess } from '../../shared/child-process/run-process'
|
||||
import type { RuntimeChildProcess } from './desktop-script-serve-channel'
|
||||
import { DesktopScriptRuntimeHost, isRuntimeHostUnavailable } from './desktop-script-runtime-host'
|
||||
|
||||
const POLICY_ERROR =
|
||||
'File runtime.ps1 cannot be loaded because running scripts is disabled on this system.\n + CategoryInfo : SecurityError'
|
||||
'File runtime.ps1 cannot be loaded because running scripts\nis disabled on this system.\n + CategoryInfo : SecurityError'
|
||||
|
||||
class FakeRuntimeChild extends EventEmitter {
|
||||
readonly stdout = new EventEmitter()
|
||||
@@ -35,33 +35,50 @@ class FakeRuntimeChild extends EventEmitter {
|
||||
return this.writes.map((line) => JSON.parse(line) as Record<string, unknown>)
|
||||
}
|
||||
|
||||
respond(response: unknown): void {
|
||||
this.stdout.emit('data', Buffer.from(`${JSON.stringify(response)}\n`, 'utf8'))
|
||||
/** The id the host is currently waiting on, so replies can echo it. */
|
||||
pendingId(): number {
|
||||
return this.requests().at(-1)?.requestId as number
|
||||
}
|
||||
|
||||
respond(response: Record<string, unknown>, requestId = this.pendingId()): void {
|
||||
this.write(`${JSON.stringify({ ...response, requestId })}\n`)
|
||||
}
|
||||
|
||||
write(raw: string): void {
|
||||
this.stdout.emit('data', Buffer.from(raw, 'utf8'))
|
||||
}
|
||||
|
||||
exit(code: number | null, stderr = ''): void {
|
||||
if (stderr) {
|
||||
this.stderr.emit('data', Buffer.from(stderr, 'utf8'))
|
||||
}
|
||||
this.emit('exit', code, null)
|
||||
this.emit('close', code, null)
|
||||
}
|
||||
}
|
||||
|
||||
function createHost(options: { idleShutdownMs?: number; requestTimeoutMs?: number } = {}) {
|
||||
function createHost(
|
||||
options: {
|
||||
idleShutdownMs?: number
|
||||
requestTimeoutMs?: number
|
||||
cooldownMs?: number
|
||||
now?: () => number
|
||||
} = {}
|
||||
) {
|
||||
const children: FakeRuntimeChild[] = []
|
||||
const specs: ProcessSpec[] = []
|
||||
const warnings: string[] = []
|
||||
const host = new DesktopScriptRuntimeHost('C:\\orca\\runtime.ps1', {
|
||||
...options,
|
||||
powerShellPath: () => 'C:\\Windows\\System32\\powershell.exe',
|
||||
warn: () => {},
|
||||
warn: (message) => warnings.push(message),
|
||||
spawn: (spec) => {
|
||||
specs.push(spec)
|
||||
const child = new FakeRuntimeChild()
|
||||
children.push(child)
|
||||
return child as unknown as ReturnType<typeof spawnProcess>
|
||||
return child as unknown as RuntimeChildProcess
|
||||
}
|
||||
})
|
||||
return { host, children, specs }
|
||||
return { host, children, specs, warnings }
|
||||
}
|
||||
|
||||
/** Let the host's queue microtasks drain so the next request reaches its child. */
|
||||
@@ -71,6 +88,17 @@ async function settle(): Promise<void> {
|
||||
}
|
||||
}
|
||||
|
||||
/** Kill each helper the host starts, until it stops starting them. */
|
||||
async function failEveryStart(children: FakeRuntimeChild[], stderr: string): Promise<void> {
|
||||
for (let index = 0; index < 8; index++) {
|
||||
if (index >= children.length) {
|
||||
return
|
||||
}
|
||||
children[index].exit(1, stderr)
|
||||
await settle()
|
||||
}
|
||||
}
|
||||
|
||||
describe('DesktopScriptRuntimeHost', () => {
|
||||
afterEach(() => {
|
||||
vi.useRealTimers()
|
||||
@@ -94,6 +122,7 @@ describe('DesktopScriptRuntimeHost', () => {
|
||||
expect(children).toHaveLength(1)
|
||||
expect(children[0].requests()).toHaveLength(6)
|
||||
expect(specs[0].args).toEqual([
|
||||
'-NoLogo',
|
||||
'-NoProfile',
|
||||
'-NonInteractive',
|
||||
'-ExecutionPolicy',
|
||||
@@ -112,7 +141,7 @@ describe('DesktopScriptRuntimeHost', () => {
|
||||
const second = host.request({ tool: 'click', app: 'B' })
|
||||
await settle()
|
||||
|
||||
expect(children[0].requests()).toEqual([{ tool: 'click', app: 'A' }])
|
||||
expect(children[0].requests()).toEqual([{ tool: 'click', app: 'A', requestId: 1 }])
|
||||
|
||||
children[0].respond({ ok: true, action: { path: 'synthetic' } })
|
||||
await expect(first).resolves.toMatchObject({ ok: true })
|
||||
@@ -124,13 +153,23 @@ describe('DesktopScriptRuntimeHost', () => {
|
||||
host.dispose()
|
||||
})
|
||||
|
||||
it('strips the echoed id from the response it hands back', async () => {
|
||||
const { host, children } = createHost()
|
||||
const promise = host.request({ tool: 'handshake' })
|
||||
await settle()
|
||||
children[0].respond({ ok: true, capabilities: {} })
|
||||
|
||||
await expect(promise).resolves.toEqual({ ok: true, capabilities: {} })
|
||||
host.dispose()
|
||||
})
|
||||
|
||||
it('reassembles a response split across chunks, including a split code point', async () => {
|
||||
const { host, children } = createHost()
|
||||
const promise = host.request({ tool: 'get_app_state', app: 'Editor' })
|
||||
await settle()
|
||||
|
||||
const payload = Buffer.from(
|
||||
`${JSON.stringify({ ok: true, snapshot: { app: 'né' } })}\n`,
|
||||
`${JSON.stringify({ ok: true, snapshot: { app: 'né' }, requestId: 1 })}\r\n`,
|
||||
'utf8'
|
||||
)
|
||||
const split = payload.indexOf(Buffer.from('é', 'utf8')) + 1
|
||||
@@ -141,6 +180,32 @@ describe('DesktopScriptRuntimeHost', () => {
|
||||
host.dispose()
|
||||
})
|
||||
|
||||
it('kills the helper rather than answering a request with another reply', async () => {
|
||||
const { host, children } = createHost()
|
||||
|
||||
const first = host.request({ tool: 'handshake' })
|
||||
await settle()
|
||||
// A stray line would otherwise shift every later response by one.
|
||||
children[0].respond({ ok: true, capabilities: {} }, 999)
|
||||
|
||||
await expect(first).rejects.toThrow(/did not match the pending request/)
|
||||
expect(children[0].killed).toBe(true)
|
||||
host.dispose()
|
||||
})
|
||||
|
||||
it('kills the helper when an unsolicited line arrives with nothing pending', async () => {
|
||||
const { host, children } = createHost()
|
||||
|
||||
const first = host.request({ tool: 'handshake' })
|
||||
await settle()
|
||||
children[0].respond({ ok: true, capabilities: {} })
|
||||
await first
|
||||
|
||||
children[0].write(`${JSON.stringify({ ok: true, requestId: 77 })}\n`)
|
||||
expect(children[0].killed).toBe(true)
|
||||
host.dispose()
|
||||
})
|
||||
|
||||
it('times out a wedged operation and starts a fresh helper for the next one', async () => {
|
||||
vi.useFakeTimers()
|
||||
const { host, children } = createHost({ requestTimeoutMs: 30_000 })
|
||||
@@ -183,6 +248,56 @@ describe('DesktopScriptRuntimeHost', () => {
|
||||
host.dispose()
|
||||
})
|
||||
|
||||
it('stops respawning a helper that dies on every second operation', async () => {
|
||||
let clock = 1_000
|
||||
const { host, children } = createHost({ cooldownMs: 60_000, now: () => clock })
|
||||
|
||||
// One good answer per helper is exactly the pattern that used to respawn
|
||||
// forever: the success reset the failure count before it could ever trip.
|
||||
for (let round = 0; round < 3; round++) {
|
||||
const good = host.request({ tool: 'handshake' })
|
||||
await settle()
|
||||
children.at(-1)?.respond({ ok: true, capabilities: {} })
|
||||
await expect(good).resolves.toMatchObject({ ok: true })
|
||||
await settle()
|
||||
|
||||
const crash = host.request({ tool: 'click', app: 'Crashy' })
|
||||
await settle()
|
||||
children.at(-1)?.exit(1, 'boom')
|
||||
await expect(crash).rejects.toThrow(/runtime host exited/)
|
||||
await settle()
|
||||
}
|
||||
|
||||
const spawned = children.length
|
||||
await expect(host.request({ tool: 'handshake' })).rejects.toSatisfy(isRuntimeHostUnavailable)
|
||||
expect(children).toHaveLength(spawned)
|
||||
host.dispose()
|
||||
})
|
||||
|
||||
it('keeps serving a healthy helper after an isolated crash', async () => {
|
||||
const { host, children } = createHost({ cooldownMs: 60_000 })
|
||||
|
||||
const crashed = host.request({ tool: 'handshake' })
|
||||
await settle()
|
||||
children[0].respond({ ok: true, capabilities: {} })
|
||||
await crashed
|
||||
const second = host.request({ tool: 'click', app: 'Notepad' })
|
||||
await settle()
|
||||
children[0].exit(1, 'boom')
|
||||
await expect(second).rejects.toThrow(/runtime host exited/)
|
||||
|
||||
for (let index = 0; index < 4; index++) {
|
||||
const next = host.request({ tool: 'handshake' })
|
||||
await settle()
|
||||
children.at(-1)?.respond({ ok: true, capabilities: {} })
|
||||
await expect(next).resolves.toMatchObject({ ok: true })
|
||||
}
|
||||
|
||||
// A clean run clears the count, so one bad helper cannot degrade a good one.
|
||||
expect(children).toHaveLength(2)
|
||||
host.dispose()
|
||||
})
|
||||
|
||||
it('shuts the helper down when idle and starts a new one on the next operation', async () => {
|
||||
vi.useFakeTimers()
|
||||
const { host, children } = createHost({ idleShutdownMs: 60_000 })
|
||||
@@ -218,8 +333,22 @@ describe('DesktopScriptRuntimeHost', () => {
|
||||
await expect(promise).rejects.toThrow(/shut down/)
|
||||
})
|
||||
|
||||
it('never respawns for a request queued behind dispose', async () => {
|
||||
const { host, children } = createHost()
|
||||
const first = host.request({ tool: 'handshake' })
|
||||
const queued = host.request({ tool: 'handshake' })
|
||||
await settle()
|
||||
|
||||
host.dispose()
|
||||
await expect(first).rejects.toBeInstanceOf(Error)
|
||||
await expect(queued).rejects.toSatisfy(isRuntimeHostUnavailable)
|
||||
await settle()
|
||||
|
||||
expect(children).toHaveLength(1)
|
||||
})
|
||||
|
||||
it('falls back to Bypass once when the execution policy blocks the start', async () => {
|
||||
const { host, children, specs } = createHost()
|
||||
const { host, children, specs, warnings } = createHost()
|
||||
|
||||
const promise = host.request({ tool: 'handshake' })
|
||||
await settle()
|
||||
@@ -230,6 +359,7 @@ describe('DesktopScriptRuntimeHost', () => {
|
||||
expect(specs[1].args).toContain('Bypass')
|
||||
children[1].respond({ ok: true, capabilities: {} })
|
||||
await expect(promise).resolves.toMatchObject({ ok: true })
|
||||
expect(warnings.some((line) => /Bypass for the rest of this session/.test(line))).toBe(true)
|
||||
|
||||
// The fallback is remembered for the session rather than re-probed per call.
|
||||
const next = host.request({ tool: 'handshake' })
|
||||
@@ -245,12 +375,10 @@ describe('DesktopScriptRuntimeHost', () => {
|
||||
|
||||
const promise = host.request({ tool: 'handshake' })
|
||||
await settle()
|
||||
children[0].exit(1, POLICY_ERROR)
|
||||
await settle()
|
||||
children[1].exit(1, POLICY_ERROR)
|
||||
await failEveryStart(children, POLICY_ERROR)
|
||||
|
||||
await expect(promise).rejects.toSatisfy(isRuntimeHostUnavailable)
|
||||
await expect(host.request({ tool: 'handshake' })).rejects.toSatisfy(isRuntimeHostUnavailable)
|
||||
host.dispose()
|
||||
})
|
||||
|
||||
it('reports itself unavailable when the helper cannot be spawned at all', async () => {
|
||||
@@ -263,15 +391,50 @@ describe('DesktopScriptRuntimeHost', () => {
|
||||
})
|
||||
|
||||
await expect(host.request({ tool: 'handshake' })).rejects.toSatisfy(isRuntimeHostUnavailable)
|
||||
host.dispose()
|
||||
})
|
||||
|
||||
it('reports itself unavailable when a fresh helper dies before answering', async () => {
|
||||
it('retries a transient pre-answer death without the caller ever seeing it', async () => {
|
||||
const { host, children } = createHost()
|
||||
|
||||
const promise = host.request({ tool: 'handshake' })
|
||||
await settle()
|
||||
children[0].exit(1, 'The term is not recognized')
|
||||
children[0].exit(1, 'Add-Type : Cannot access the temporary directory')
|
||||
await settle()
|
||||
|
||||
await expect(promise).rejects.toSatisfy(isRuntimeHostUnavailable)
|
||||
expect(children).toHaveLength(2)
|
||||
children[1].respond({ ok: true, capabilities: {} })
|
||||
|
||||
await expect(promise).resolves.toMatchObject({ ok: true })
|
||||
host.dispose()
|
||||
})
|
||||
|
||||
it('gives up only after repeated start failures, then serves from the host again after the cooldown', async () => {
|
||||
let clock = 1_000
|
||||
const { host, children, warnings } = createHost({ cooldownMs: 60_000, now: () => clock })
|
||||
|
||||
const failed = host.request({ tool: 'handshake' })
|
||||
await settle()
|
||||
await failEveryStart(children, 'The term is not recognized')
|
||||
await expect(failed).rejects.toSatisfy(isRuntimeHostUnavailable)
|
||||
|
||||
const attempts = children.length
|
||||
expect(attempts).toBe(3)
|
||||
|
||||
// Inside the cooldown the host stays out of the way without respawning.
|
||||
clock += 30_000
|
||||
await expect(host.request({ tool: 'handshake' })).rejects.toSatisfy(isRuntimeHostUnavailable)
|
||||
expect(children).toHaveLength(attempts)
|
||||
|
||||
// Past it, the next operation re-probes rather than staying degraded forever.
|
||||
clock += 31_000
|
||||
const recovered = host.request({ tool: 'handshake' })
|
||||
await settle()
|
||||
expect(children).toHaveLength(attempts + 1)
|
||||
children[attempts].respond({ ok: true, capabilities: {} })
|
||||
await expect(recovered).resolves.toMatchObject({ ok: true })
|
||||
|
||||
expect(warnings.at(-1)).toMatch(/recovered/)
|
||||
host.dispose()
|
||||
})
|
||||
})
|
||||
|
||||
@@ -1,36 +1,41 @@
|
||||
import { StringDecoder } from 'node:string_decoder'
|
||||
import type { ProcessSpec } from '../../shared/child-process/process-spec'
|
||||
import { spawnProcess } from '../../shared/child-process/run-process'
|
||||
import { windowsPowerShellPath } from '../../shared/child-process/windows-system-binary'
|
||||
import { reportComputerDiagnostic } from './computer-sidecar-diagnostics'
|
||||
import type { BridgeRequest, BridgeResponse } from './desktop-script-provider-types'
|
||||
import {
|
||||
startServeChannel,
|
||||
type DesktopScriptServeChannel,
|
||||
type RuntimeProcessSpawn
|
||||
} from './desktop-script-serve-channel'
|
||||
import {
|
||||
MAX_START_ATTEMPTS,
|
||||
RuntimeHostAvailability,
|
||||
START_FAILURE_COOLDOWN_MS
|
||||
} from './desktop-script-runtime-availability'
|
||||
import { RuntimeClientError } from './runtime-client-error'
|
||||
import {
|
||||
FALLBACK_WINDOWS_EXECUTION_POLICY,
|
||||
PREFERRED_WINDOWS_EXECUTION_POLICY,
|
||||
isExecutionPolicyBlocked,
|
||||
windowsPowerShellRuntimeArgs,
|
||||
type WindowsExecutionPolicy
|
||||
windowsPowerShellRuntimeArgs
|
||||
} from './windows-powershell-execution-policy'
|
||||
|
||||
/** The all-pipes child `spawnProcess` returns; avoids a node:child_process import. */
|
||||
type RuntimeChildProcess = ReturnType<typeof spawnProcess>
|
||||
|
||||
const REQUEST_TIMEOUT_MS = 30_000
|
||||
const IDLE_SHUTDOWN_MS = 120_000
|
||||
const MAX_RESPONSE_BYTES = 20 * 1024 * 1024
|
||||
|
||||
/** Code the client keys on to fall back to the one-shot bridge for the session. */
|
||||
/** Code the client keys on to serve this one operation from the one-shot bridge. */
|
||||
export const RUNTIME_HOST_UNAVAILABLE = 'runtime_host_unavailable'
|
||||
|
||||
export type DesktopScriptRuntimeHostOptions = {
|
||||
spawn?: (spec: ProcessSpec) => RuntimeChildProcess
|
||||
spawn?: RuntimeProcessSpawn
|
||||
powerShellPath?: () => string
|
||||
requestTimeoutMs?: number
|
||||
idleShutdownMs?: number
|
||||
cooldownMs?: number
|
||||
now?: () => number
|
||||
warn?: (message: string) => void
|
||||
}
|
||||
|
||||
type PendingRequest = {
|
||||
id: number
|
||||
resolve: (response: BridgeResponse) => void
|
||||
reject: (error: Error) => void
|
||||
timer: NodeJS.Timeout
|
||||
@@ -50,23 +55,21 @@ export function isRuntimeHostUnavailable(error: unknown): boolean {
|
||||
* screen capture. Compiling once per session collapses a burst of short-lived
|
||||
* PIDs into a single process.
|
||||
*
|
||||
* Requests are strictly serialized: the protocol carries no request id because
|
||||
* only one operation is ever in flight, and native automation is not safe to
|
||||
* interleave anyway.
|
||||
* Requests are strictly serialized, and each carries an id the helper echoes.
|
||||
* Serialization alone would leave a single stray line answering every later
|
||||
* request with the previous response — silently acting on stale element
|
||||
* indexes, with no error raised — so the id is checked and a mismatch is fatal
|
||||
* to the child rather than merely logged.
|
||||
*/
|
||||
export class DesktopScriptRuntimeHost {
|
||||
private child: RuntimeChildProcess | null = null
|
||||
private detachChild: (() => void) | null = null
|
||||
private decoder = new StringDecoder('utf8')
|
||||
private stdoutBuffer = ''
|
||||
private stderrText = ''
|
||||
private channel: DesktopScriptServeChannel | null = null
|
||||
private pending: PendingRequest | null = null
|
||||
private queueTail: Promise<void> | null = null
|
||||
private idleTimer: NodeJS.Timeout | null = null
|
||||
private policy: WindowsExecutionPolicy = PREFERRED_WINDOWS_EXECUTION_POLICY
|
||||
private policyRetryPending = false
|
||||
private childAnswered = false
|
||||
private unavailable = false
|
||||
private disposed = false
|
||||
private nextRequestId = 1
|
||||
private readonly availability: RuntimeHostAvailability
|
||||
private readonly requestTimeoutMs: number
|
||||
private readonly idleShutdownMs: number
|
||||
|
||||
@@ -76,6 +79,11 @@ export class DesktopScriptRuntimeHost {
|
||||
) {
|
||||
this.requestTimeoutMs = options.requestTimeoutMs ?? REQUEST_TIMEOUT_MS
|
||||
this.idleShutdownMs = options.idleShutdownMs ?? IDLE_SHUTDOWN_MS
|
||||
this.availability = new RuntimeHostAvailability(
|
||||
options.cooldownMs ?? START_FAILURE_COOLDOWN_MS,
|
||||
options.now ?? Date.now,
|
||||
(message) => (options.warn ?? reportComputerDiagnostic)(message)
|
||||
)
|
||||
}
|
||||
|
||||
request(request: BridgeRequest): Promise<BridgeResponse> {
|
||||
@@ -96,10 +104,12 @@ export class DesktopScriptRuntimeHost {
|
||||
return result
|
||||
}
|
||||
|
||||
/** Stop the helper. A later request starts a fresh one. */
|
||||
/** Permanently stop this host. Callers build a new one for a new session. */
|
||||
dispose(): void {
|
||||
this.disposed = true
|
||||
this.clearIdleTimer()
|
||||
this.stopChild()
|
||||
this.availability.clearCooldown()
|
||||
this.stopChannel()
|
||||
this.rejectPending(
|
||||
new RuntimeClientError('accessibility_error', 'desktop provider runtime host was shut down')
|
||||
)
|
||||
@@ -107,39 +117,59 @@ export class DesktopScriptRuntimeHost {
|
||||
|
||||
private async send(request: BridgeRequest): Promise<BridgeResponse> {
|
||||
this.clearIdleTimer()
|
||||
try {
|
||||
return await this.sendOnce(request)
|
||||
} catch (error) {
|
||||
if (!this.policyRetryPending) {
|
||||
throw error
|
||||
}
|
||||
this.policyRetryPending = false
|
||||
this.policy = FALLBACK_WINDOWS_EXECUTION_POLICY
|
||||
this.warn(
|
||||
`runtime host start blocked at ${PREFERRED_WINDOWS_EXECUTION_POLICY}; retrying once with ${FALLBACK_WINDOWS_EXECUTION_POLICY}`
|
||||
)
|
||||
return await this.sendOnce(request)
|
||||
// Why checked here and not only on entry: requests queue, and dispose can
|
||||
// land while one waits its turn. Without this a teardown respawns a helper.
|
||||
if (this.disposed) {
|
||||
throw this.unavailableError('runtime host was disposed')
|
||||
}
|
||||
const cooldown = this.availability.remainingCooldown()
|
||||
if (cooldown > 0) {
|
||||
throw this.unavailableError(`retrying the runtime host in ${cooldown}ms`)
|
||||
}
|
||||
let lastError: unknown
|
||||
for (let attempt = 1; attempt <= MAX_START_ATTEMPTS; attempt++) {
|
||||
try {
|
||||
const response = await this.sendOnce(request)
|
||||
this.availability.recordSuccess()
|
||||
return response
|
||||
} catch (error) {
|
||||
lastError = error
|
||||
if (this.availability.policyRetryPending) {
|
||||
this.availability.escalateExecutionPolicy()
|
||||
continue
|
||||
}
|
||||
// A helper that answered and then died is a crash, not a bad start: the
|
||||
// caller sees it and the next operation gets a fresh process — unless it
|
||||
// keeps happening, which is thrash the one-shot bridge should absorb.
|
||||
if (!isRuntimeHostUnavailable(error)) {
|
||||
if (this.availability.exhausted) {
|
||||
this.availability.enterCooldown()
|
||||
}
|
||||
throw error
|
||||
}
|
||||
this.availability.warn(
|
||||
`runtime host failed to start (attempt ${attempt}/${MAX_START_ATTEMPTS}): ${errorText(error)}`
|
||||
)
|
||||
}
|
||||
}
|
||||
this.availability.enterCooldown()
|
||||
throw lastError
|
||||
}
|
||||
|
||||
private sendOnce(request: BridgeRequest): Promise<BridgeResponse> {
|
||||
if (this.unavailable) {
|
||||
return Promise.reject(this.unavailableError('runtime host is unavailable'))
|
||||
}
|
||||
let child: RuntimeChildProcess
|
||||
let channel: DesktopScriptServeChannel
|
||||
try {
|
||||
child = this.ensureChild()
|
||||
channel = this.ensureChannel()
|
||||
} catch (error) {
|
||||
this.unavailable = true
|
||||
return Promise.reject(
|
||||
this.unavailableError(error instanceof Error ? error.message : String(error))
|
||||
)
|
||||
this.availability.recordFailure()
|
||||
return Promise.reject(this.unavailableError(errorText(error)))
|
||||
}
|
||||
const id = this.nextRequestId++
|
||||
return new Promise((resolve, reject) => {
|
||||
// Why kill rather than wait: a hung UI Automation call cannot be
|
||||
// cancelled, so the process itself is the only thing left to reclaim.
|
||||
const timer = setTimeout(() => {
|
||||
this.abortChild(
|
||||
this.abortChannel(
|
||||
new RuntimeClientError(
|
||||
'action_timeout',
|
||||
`desktop provider timed out after ${this.requestTimeoutMs}ms`
|
||||
@@ -147,145 +177,107 @@ export class DesktopScriptRuntimeHost {
|
||||
)
|
||||
}, this.requestTimeoutMs)
|
||||
timer.unref?.()
|
||||
this.pending = { resolve, reject, timer }
|
||||
child.stdin.write(`${JSON.stringify(request)}\n`, (error) => {
|
||||
if (error) {
|
||||
this.abortChild(new RuntimeClientError('accessibility_error', error.message))
|
||||
}
|
||||
})
|
||||
})
|
||||
}
|
||||
|
||||
private ensureChild(): RuntimeChildProcess {
|
||||
if (this.child) {
|
||||
return this.child
|
||||
}
|
||||
const spawn = this.options.spawn ?? spawnProcess
|
||||
const child = spawn({
|
||||
program: (this.options.powerShellPath ?? windowsPowerShellPath)(),
|
||||
args: windowsPowerShellRuntimeArgs(this.scriptPath, this.policy, ['-Serve']),
|
||||
env: process.env
|
||||
})
|
||||
this.decoder = new StringDecoder('utf8')
|
||||
this.stdoutBuffer = ''
|
||||
this.stderrText = ''
|
||||
this.childAnswered = false
|
||||
|
||||
const onStdout = (chunk: Buffer | string): void => this.readStdout(chunk)
|
||||
const onStderr = (chunk: Buffer | string): void => {
|
||||
this.stderrText = `${this.stderrText}${chunk.toString()}`.slice(-4096)
|
||||
}
|
||||
const onExit = (code: number | null, signal: NodeJS.Signals | null): void =>
|
||||
this.handleGone(child, signal ? `signal ${signal}` : `code ${code ?? 'unknown'}`)
|
||||
const onError = (error: Error): void => this.handleGone(child, error.message)
|
||||
child.stdout.on('data', onStdout)
|
||||
child.stderr.on('data', onStderr)
|
||||
child.once('exit', onExit)
|
||||
child.once('error', onError)
|
||||
// An unhandled stream error is an uncaught exception in the main process.
|
||||
child.stdin.on('error', () => {})
|
||||
this.detachChild = (): void => {
|
||||
child.stdout.off('data', onStdout)
|
||||
child.stderr.off('data', onStderr)
|
||||
child.off('exit', onExit)
|
||||
child.off('error', onError)
|
||||
child.on('error', () => {})
|
||||
}
|
||||
this.child = child
|
||||
return child
|
||||
}
|
||||
|
||||
private readStdout(chunk: Buffer | string): void {
|
||||
this.stdoutBuffer += typeof chunk === 'string' ? chunk : this.decoder.write(chunk)
|
||||
if (this.stdoutBuffer.length > MAX_RESPONSE_BYTES) {
|
||||
this.abortChild(
|
||||
new RuntimeClientError(
|
||||
'accessibility_error',
|
||||
'desktop provider response exceeded the runtime host buffer'
|
||||
)
|
||||
this.pending = { id, resolve, reject, timer }
|
||||
channel.write(`${JSON.stringify({ ...request, requestId: id })}\n`, (error) =>
|
||||
this.abortChannel(new RuntimeClientError('accessibility_error', error.message))
|
||||
)
|
||||
return
|
||||
})
|
||||
}
|
||||
|
||||
private ensureChannel(): DesktopScriptServeChannel {
|
||||
if (this.channel) {
|
||||
return this.channel
|
||||
}
|
||||
for (let newline = this.stdoutBuffer.indexOf('\n'); newline >= 0;) {
|
||||
const line = this.stdoutBuffer.slice(0, newline).trim()
|
||||
this.stdoutBuffer = this.stdoutBuffer.slice(newline + 1)
|
||||
if (line) {
|
||||
this.deliver(line)
|
||||
this.childAnswered = false
|
||||
const channel: DesktopScriptServeChannel = startServeChannel(
|
||||
{
|
||||
program: (this.options.powerShellPath ?? windowsPowerShellPath)(),
|
||||
args: windowsPowerShellRuntimeArgs(this.scriptPath, this.availability.executionPolicy, [
|
||||
'-Serve'
|
||||
]),
|
||||
env: process.env
|
||||
},
|
||||
this.options.spawn ?? spawnProcess,
|
||||
{
|
||||
onLine: (line) => this.deliver(line),
|
||||
// A replaced channel can still report; that must not fail the live one.
|
||||
onGone: (detail) => {
|
||||
if (this.channel === channel) {
|
||||
this.handleGone(detail)
|
||||
}
|
||||
},
|
||||
onOverflow: () =>
|
||||
this.abortChannel(
|
||||
new RuntimeClientError(
|
||||
'accessibility_error',
|
||||
'desktop provider response exceeded the runtime host buffer'
|
||||
)
|
||||
)
|
||||
}
|
||||
newline = this.stdoutBuffer.indexOf('\n')
|
||||
}
|
||||
)
|
||||
this.channel = channel
|
||||
return channel
|
||||
}
|
||||
|
||||
private deliver(line: string): void {
|
||||
this.childAnswered = true
|
||||
const pending = this.takePending()
|
||||
if (!pending) {
|
||||
let parsed: BridgeResponse
|
||||
try {
|
||||
parsed = JSON.parse(line) as BridgeResponse
|
||||
} catch {
|
||||
// Not a response at all — a PowerShell banner, a stray write. Dropping it
|
||||
// is safe now that the id below is what decides which request is answered,
|
||||
// and it keeps a chatty console from making the helper unusable.
|
||||
return
|
||||
}
|
||||
try {
|
||||
pending.resolve(JSON.parse(line) as BridgeResponse)
|
||||
} catch (error) {
|
||||
pending.reject(
|
||||
const pending = this.pending
|
||||
if (!pending || parsed.requestId !== pending.id) {
|
||||
// One unmatched reply would otherwise shift every later response by one.
|
||||
this.abortChannel(
|
||||
new RuntimeClientError(
|
||||
'accessibility_error',
|
||||
`desktop provider returned invalid JSON: ${error instanceof Error ? error.message : String(error)}`
|
||||
'desktop provider response did not match the pending request'
|
||||
)
|
||||
)
|
||||
}
|
||||
}
|
||||
|
||||
private handleGone(child: RuntimeChildProcess, detail: string): void {
|
||||
// A replaced child can still report; that must not fail the live one.
|
||||
if (this.child !== child) {
|
||||
return
|
||||
}
|
||||
const text = [detail, this.stderrText.trim()].filter(Boolean).join(': ')
|
||||
// Only a reply this host can prove is its own counts as the helper working.
|
||||
this.childAnswered = true
|
||||
this.pending = null
|
||||
clearTimeout(pending.timer)
|
||||
const { requestId: _echoed, ...response } = parsed
|
||||
pending.resolve(response)
|
||||
}
|
||||
|
||||
private handleGone(detail: string): void {
|
||||
const answered = this.childAnswered
|
||||
this.releaseChild()
|
||||
if (
|
||||
!answered &&
|
||||
this.policy === PREFERRED_WINDOWS_EXECUTION_POLICY &&
|
||||
isExecutionPolicyBlocked(text)
|
||||
) {
|
||||
this.policyRetryPending = true
|
||||
this.rejectPending(new RuntimeClientError('accessibility_error', text))
|
||||
this.channel = null
|
||||
this.availability.recordFailure()
|
||||
if (!answered && this.availability.atPreferredPolicy && isExecutionPolicyBlocked(detail)) {
|
||||
this.availability.requestPolicyRetry()
|
||||
this.rejectPending(new RuntimeClientError('accessibility_error', detail))
|
||||
return
|
||||
}
|
||||
if (!answered) {
|
||||
this.unavailable = true
|
||||
this.rejectPending(this.unavailableError(text))
|
||||
this.rejectPending(this.unavailableError(detail))
|
||||
return
|
||||
}
|
||||
this.rejectPending(
|
||||
new RuntimeClientError('accessibility_error', `desktop provider runtime host exited: ${text}`)
|
||||
new RuntimeClientError(
|
||||
'accessibility_error',
|
||||
`desktop provider runtime host exited: ${detail}`
|
||||
)
|
||||
)
|
||||
}
|
||||
|
||||
private abortChild(error: Error): void {
|
||||
this.stopChild()
|
||||
private abortChannel(error: Error): void {
|
||||
this.stopChannel()
|
||||
this.rejectPending(error)
|
||||
}
|
||||
|
||||
private stopChild(): void {
|
||||
const child = this.child
|
||||
this.releaseChild()
|
||||
if (!child) {
|
||||
return
|
||||
}
|
||||
// Closing stdin ends the serve loop; the kill covers a wedged helper.
|
||||
try {
|
||||
child.stdin.end()
|
||||
} catch {
|
||||
/* already closed */
|
||||
}
|
||||
child.kill()
|
||||
}
|
||||
|
||||
private releaseChild(): void {
|
||||
const detach = this.detachChild
|
||||
this.detachChild = null
|
||||
detach?.()
|
||||
this.child = null
|
||||
private stopChannel(): void {
|
||||
const channel = this.channel
|
||||
this.channel = null
|
||||
channel?.stop()
|
||||
}
|
||||
|
||||
private takePending(): PendingRequest | null {
|
||||
@@ -303,12 +295,12 @@ export class DesktopScriptRuntimeHost {
|
||||
|
||||
private armIdleTimer(): void {
|
||||
this.clearIdleTimer()
|
||||
if (!this.child) {
|
||||
if (!this.channel) {
|
||||
return
|
||||
}
|
||||
this.idleTimer = setTimeout(() => {
|
||||
this.idleTimer = null
|
||||
this.stopChild()
|
||||
this.stopChannel()
|
||||
}, this.idleShutdownMs)
|
||||
this.idleTimer.unref?.()
|
||||
}
|
||||
@@ -326,8 +318,8 @@ export class DesktopScriptRuntimeHost {
|
||||
`desktop provider runtime host could not start: ${message}`
|
||||
)
|
||||
}
|
||||
}
|
||||
|
||||
private warn(message: string): void {
|
||||
;(this.options.warn ?? ((text: string) => console.warn(`[computer-use] ${text}`)))(message)
|
||||
}
|
||||
function errorText(error: unknown): string {
|
||||
return error instanceof Error ? error.message : String(error)
|
||||
}
|
||||
|
||||
@@ -0,0 +1,138 @@
|
||||
import { StringDecoder } from 'node:string_decoder'
|
||||
import type { ProcessSpec } from '../../shared/child-process/process-spec'
|
||||
import type { spawnProcess } from '../../shared/child-process/run-process'
|
||||
|
||||
/** The all-pipes child `spawnProcess` returns; avoids a node:child_process import. */
|
||||
export type RuntimeChildProcess = ReturnType<typeof spawnProcess>
|
||||
|
||||
export type RuntimeProcessSpawn = (spec: ProcessSpec) => RuntimeChildProcess
|
||||
|
||||
/** UTF-16 units, not bytes — this bounds the buffer, it is not a payload contract. */
|
||||
const MAX_RESPONSE_CHARS = 20 * 1024 * 1024
|
||||
const MAX_STDERR_CHARS = 4096
|
||||
|
||||
export type ServeChannelHandlers = {
|
||||
/** One complete line from the helper, without its terminator. */
|
||||
onLine: (line: string) => void
|
||||
/** The helper is gone; detail carries the exit reason and its stderr tail. */
|
||||
onGone: (detail: string) => void
|
||||
/** The helper produced more than one buffer's worth without a line break. */
|
||||
onOverflow: () => void
|
||||
}
|
||||
|
||||
/**
|
||||
* One `runtime.ps1 -Serve` child, framed as NDJSON lines.
|
||||
*
|
||||
* Split from the host so the host reads as what it is — a queue, a retry policy
|
||||
* and a correlation check — rather than that plus stream plumbing. Responses
|
||||
* carry base64 screenshots and routinely exceed a megabyte, so lines are
|
||||
* reassembled across chunks with a decoder that survives a code point split
|
||||
* across a chunk boundary.
|
||||
*/
|
||||
export class DesktopScriptServeChannel {
|
||||
private readonly decoder = new StringDecoder('utf8')
|
||||
private buffer = ''
|
||||
private stderrTail = ''
|
||||
private detach: (() => void) | null = null
|
||||
private closed = false
|
||||
|
||||
constructor(
|
||||
private readonly child: RuntimeChildProcess,
|
||||
private readonly handlers: ServeChannelHandlers
|
||||
) {
|
||||
const onStdout = (chunk: Buffer | string): void => this.readStdout(chunk)
|
||||
const onStderr = (chunk: Buffer | string): void => {
|
||||
this.stderrTail = `${this.stderrTail}${chunk.toString()}`.slice(-MAX_STDERR_CHARS)
|
||||
}
|
||||
// Why close and not exit: the caller classifies the failure from stderr, and
|
||||
// only close guarantees the stdio streams were drained first.
|
||||
const onClose = (code: number | null, signal: NodeJS.Signals | null): void =>
|
||||
this.reportGone(signal ? `signal ${signal}` : `code ${code ?? 'unknown'}`)
|
||||
const onError = (error: Error): void => this.reportGone(error.message)
|
||||
child.stdout.on('data', onStdout)
|
||||
child.stderr.on('data', onStderr)
|
||||
child.once('close', onClose)
|
||||
child.once('error', onError)
|
||||
// An unhandled stream error is an uncaught exception in the main process.
|
||||
child.stdin.on('error', () => {})
|
||||
this.detach = (): void => {
|
||||
child.stdout.off('data', onStdout)
|
||||
child.stderr.off('data', onStderr)
|
||||
child.off('close', onClose)
|
||||
child.off('error', onError)
|
||||
child.on('error', () => {})
|
||||
}
|
||||
}
|
||||
|
||||
write(payload: string, onError: (error: Error) => void): void {
|
||||
this.child.stdin.write(payload, (error) => {
|
||||
if (error) {
|
||||
onError(error)
|
||||
}
|
||||
})
|
||||
}
|
||||
|
||||
/** Stop the helper and go silent; handlers are not called afterwards. */
|
||||
stop(): void {
|
||||
if (this.closed) {
|
||||
return
|
||||
}
|
||||
this.closed = true
|
||||
this.detach?.()
|
||||
this.detach = null
|
||||
this.buffer = ''
|
||||
// Closing stdin ends the serve loop; the kill covers a wedged helper.
|
||||
try {
|
||||
this.child.stdin.end()
|
||||
} catch {
|
||||
/* already closed */
|
||||
}
|
||||
this.child.kill()
|
||||
}
|
||||
|
||||
private reportGone(detail: string): void {
|
||||
if (this.closed) {
|
||||
return
|
||||
}
|
||||
const text = [detail, this.stderrTail.trim()].filter(Boolean).join(': ')
|
||||
this.closed = true
|
||||
this.detach?.()
|
||||
this.detach = null
|
||||
this.handlers.onGone(text)
|
||||
}
|
||||
|
||||
private readStdout(chunk: Buffer | string): void {
|
||||
if (this.closed) {
|
||||
return
|
||||
}
|
||||
this.buffer += typeof chunk === 'string' ? chunk : this.decoder.write(chunk)
|
||||
if (this.buffer.length > MAX_RESPONSE_CHARS) {
|
||||
this.buffer = ''
|
||||
this.handlers.onOverflow()
|
||||
return
|
||||
}
|
||||
for (let newline = this.buffer.indexOf('\n'); newline >= 0;) {
|
||||
// Slice a trailing CR off by index; trimming copies the whole payload.
|
||||
const end = newline > 0 && this.buffer.charCodeAt(newline - 1) === 13 ? newline - 1 : newline
|
||||
const line = this.buffer.slice(0, end)
|
||||
this.buffer = this.buffer.slice(newline + 1)
|
||||
if (line.length > 0) {
|
||||
this.handlers.onLine(line)
|
||||
// A handler may have stopped this channel; stop reading its backlog.
|
||||
if (this.closed) {
|
||||
this.buffer = ''
|
||||
return
|
||||
}
|
||||
}
|
||||
newline = this.buffer.indexOf('\n')
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
export function startServeChannel(
|
||||
spec: ProcessSpec,
|
||||
spawn: RuntimeProcessSpawn,
|
||||
handlers: ServeChannelHandlers
|
||||
): DesktopScriptServeChannel {
|
||||
return new DesktopScriptServeChannel(spawn(spec), handlers)
|
||||
}
|
||||
@@ -9,6 +9,7 @@ import type {
|
||||
ComputerSnapshotResult
|
||||
} from '../../shared/runtime-types'
|
||||
import { normalizeComputerActionResult } from './computer-action-verification-normalization'
|
||||
import { isComputerSidecarDiagnostic, logComputerDiagnostic } from './computer-sidecar-diagnostics'
|
||||
import { validateComputerSidecarPasteText } from './computer-sidecar-paste-validation'
|
||||
import { RuntimeClientError } from './runtime-client-error'
|
||||
|
||||
@@ -245,6 +246,11 @@ class ComputerSidecarProcess {
|
||||
}
|
||||
|
||||
private handleMessage(message: unknown): void {
|
||||
// The sidecar's stdio is piped and unread, so its warnings arrive here.
|
||||
if (isComputerSidecarDiagnostic(message)) {
|
||||
logComputerDiagnostic(message.message)
|
||||
return
|
||||
}
|
||||
if (!isSidecarResponse(message)) {
|
||||
return
|
||||
}
|
||||
|
||||
@@ -8,6 +8,11 @@ type SidecarRequest = {
|
||||
params?: Record<string, unknown>
|
||||
}
|
||||
|
||||
// Why disconnect carries the weight on Windows: the parent stops the sidecar
|
||||
// with kill('SIGTERM'), which is TerminateProcess there, so the SIGTERM handler
|
||||
// below never runs and teardown rides on the IPC channel closing instead. A
|
||||
// helper wedged inside a UI Automation call can still outlive that and deliver
|
||||
// input after teardown; only a real signal would preempt it.
|
||||
process.once('disconnect', shutdownProviders)
|
||||
process.once('SIGTERM', () => {
|
||||
shutdownProviders()
|
||||
|
||||
@@ -13,9 +13,16 @@ export type WindowsExecutionPolicy = 'RemoteSigned' | 'Bypass'
|
||||
export const PREFERRED_WINDOWS_EXECUTION_POLICY: WindowsExecutionPolicy = 'RemoteSigned'
|
||||
export const FALLBACK_WINDOWS_EXECUTION_POLICY: WindowsExecutionPolicy = 'Bypass'
|
||||
|
||||
// Matches the SecurityError PowerShell emits for `-File` under a blocking policy.
|
||||
/**
|
||||
* Matches the SecurityError PowerShell emits for `-File` under a blocking policy.
|
||||
*
|
||||
* The prose alternative is whitespace-tolerant because PowerShell hard-wraps
|
||||
* error text at the console width, so the sentence arrives split across lines.
|
||||
* The single-token alternatives survive that wrapping unaided and are what
|
||||
* actually carries the match in practice.
|
||||
*/
|
||||
const EXECUTION_POLICY_BLOCKED =
|
||||
/running scripts is disabled on this system|UnauthorizedAccess|PSSecurityException|SecurityError|about_Execution_Policies/i
|
||||
/running\s+scripts\s+is\s+disabled|UnauthorizedAccess|PSSecurityException|about_Execution_Policies/i
|
||||
|
||||
export function isExecutionPolicyBlocked(text: string): boolean {
|
||||
return EXECUTION_POLICY_BLOCKED.test(text)
|
||||
@@ -27,6 +34,8 @@ export function windowsPowerShellRuntimeArgs(
|
||||
scriptArgs: readonly string[] = []
|
||||
): string[] {
|
||||
return [
|
||||
// -NoLogo: a banner on stdout would be read as a malformed response line.
|
||||
'-NoLogo',
|
||||
'-NoProfile',
|
||||
'-NonInteractive',
|
||||
'-ExecutionPolicy',
|
||||
|
||||
Reference in New Issue
Block a user