mirror of
https://github.com/stablyai/orca.git
synced 2026-09-29 08:03:20 +00:00
Revert "revert(ssh): take the SSH liveness readback out of this PR"
This reverts commit f0f8858964.
This commit is contained in:
@@ -0,0 +1,92 @@
|
||||
# PR #19347 round 2 — SSH exit certification moved to its owner
|
||||
|
||||
## The defect
|
||||
|
||||
`probeSshPtyLiveness` inferred "absent from `pty.listProcesses` and minted by this relay's
|
||||
generation, therefore exited". The implication is false. The relay has removal paths that never
|
||||
observe the process end: `disposePtyForRelayShutdown` waits `IMMEDIATE_PTY_EXIT_TIMEOUT_MS` for a
|
||||
physical exit, catches the timeout, and deletes the record anyway; `disposePtys` settles with
|
||||
`Promise.allSettled`, so another PTY's failed kill rejects the aggregate while that record is
|
||||
already gone; `relay-grace-lifecycle` then defers shutdown and the relay keeps serving under the
|
||||
same mint epoch. A newly connected provider asked, was told `false`, and recovery queued a live
|
||||
worker as exited.
|
||||
|
||||
Reproduced at `1e89522a9d` with the reviewer's test: 2 failed / 2 passed.
|
||||
|
||||
## Census of record removals in `PtyHandler`
|
||||
|
||||
Every removal funnels through `removePty`, which has three callers plus the shared reap.
|
||||
|
||||
| Site | Path | Observed the process end? |
|
||||
| --- | --- | --- |
|
||||
| node-pty `onExit` handler | the process reported its own exit code | Yes |
|
||||
| `reapPtyProvenExited` | pid probed, answered ESRCH from the owning host | Yes |
|
||||
| `reapExitedPty(_, 'record-torn-down')` from `listProcesses` | tidies a record already disposed elsewhere | No |
|
||||
| `disposePtyForRelayShutdown` | SIGKILL sent, exit wait timed out, record deleted regardless | No |
|
||||
|
||||
The two "yes" rows now write the id into `observedPtyExitIds`. Nothing else does, and the write
|
||||
sits at the line that made the observation rather than inside the shared reap, so the reap's
|
||||
bookkeeping caller cannot inherit it.
|
||||
|
||||
## What the relay answers
|
||||
|
||||
`pty.probeLiveness` takes an exact id and answers from the handler's own state.
|
||||
|
||||
- A live record: `live`, or `exited` if probing its pid reaps it.
|
||||
- A record mid-teardown, or an id with a revive in flight: `unknown`.
|
||||
- Any other id: `exited` only if it is in the ledger, otherwise `unknown`.
|
||||
|
||||
Fail-closed by construction. A future removal path that forgets to record itself yields `unknown`,
|
||||
which defers; it cannot yield a false death certificate. Eviction from the bounded ledger has the
|
||||
same safe direction.
|
||||
|
||||
Two gates were considered and dropped as dead code, each confirmed by a surviving mutation that
|
||||
changed no observable behaviour. A mint-epoch comparison is subsumed because ledger keys are whole
|
||||
ids that already embed the generation. A blanket "a shutdown has been fenced" refusal masked the
|
||||
ledger in the only path where teardown runs, made that path's correctness untestable, and withheld
|
||||
exits the relay really did observe.
|
||||
|
||||
## Boundary
|
||||
|
||||
**Invariant.** Only the execution host can distinguish a process that ended from bookkeeping that
|
||||
stopped tracking it. A mint epoch proves who allocated an id; it says nothing about why that
|
||||
owner's record disappeared.
|
||||
|
||||
**Owner.** `PtyHandler`. It holds the records, the physical-exit path, the pid probe, the shutdown
|
||||
fence, and pending revives.
|
||||
|
||||
**Restores or compensates.** Restores. The previous version compensated at the client, rebuilding a
|
||||
verdict from two projections that do not carry the missing evidence. The client now asks and
|
||||
forwards; `SshPtyProvider.probePtyLiveness` maps three statuses to `true` / `false` / `null` and
|
||||
does nothing else.
|
||||
|
||||
**Alternative rejected.** Keeping the inference and excluding the shutdown case client-side. It
|
||||
needs the client to know which removals were bookkeeping, which is precisely the knowledge only the
|
||||
owner has, and the next removal path would silently reopen the hole.
|
||||
|
||||
**Compatibility.** A relay too old to know the method answers JSON-RPC `-32601`, which arrives as a
|
||||
rejection and maps to `null`, so no capability negotiation is needed and no discovery round trip is
|
||||
spent. Old provider against a new relay is unaffected: the handler is additive. The `ptyIdMintEpoch`
|
||||
capability and the shared id module stay because the relay still mints through them.
|
||||
|
||||
## Mutation results
|
||||
|
||||
Eight mutations against the relay verdict and the two suites that drive it.
|
||||
|
||||
| Mutation | Result |
|
||||
| --- | --- |
|
||||
| Certify any id absent from the map | 9 failed |
|
||||
| Record the shutdown teardown as an observed exit | 4 failed |
|
||||
| Drop the node-pty `onExit` record | 2 failed |
|
||||
| Drop the pid-probe record | 1 failed |
|
||||
| Answer live instead of reaping a dead pid | 1 failed |
|
||||
| Record the torn-down sweep as an observed exit | survives |
|
||||
| Treat a disposed record as live | survives |
|
||||
| Drop the pending-revive gate | survives |
|
||||
|
||||
The three survivors guard states no current path constructs, all in the safe direction. A disposed
|
||||
record is never left in the map: all three `disposeManagedPty` callers remove it in the same
|
||||
synchronous block, so neither the torn-down sweep nor the disposed branch can be driven. An id in
|
||||
the ledger cannot be revived, because revive requires `process.kill(pid, 0)` to succeed and a
|
||||
recorded id's process is gone. The revive gate is kept anyway because that coupling lives in
|
||||
another function.
|
||||
@@ -0,0 +1,50 @@
|
||||
import { describe, expect, it, vi } from 'vitest'
|
||||
import { probeSshPtyLiveness } from './ssh-pty-liveness-probe'
|
||||
|
||||
const PTY_ID = 'pty2:0f8f3a1e-1111-4111-8111-111111111111:7'
|
||||
|
||||
function relay(answer: Error | Record<string, unknown> | null) {
|
||||
return vi.fn(async (method: string, params?: Record<string, unknown>) => {
|
||||
expect(method).toBe('pty.probeLiveness')
|
||||
expect(params).toEqual({ id: PTY_ID })
|
||||
if (answer instanceof Error) {
|
||||
throw answer
|
||||
}
|
||||
return answer
|
||||
})
|
||||
}
|
||||
|
||||
describe('SSH PTY liveness forwarder', () => {
|
||||
it.each([
|
||||
['live', true],
|
||||
['exited', false],
|
||||
['unknown', null]
|
||||
])('maps the owner verdict %s', async (status, expected) => {
|
||||
const request = relay({ status })
|
||||
|
||||
await expect(probeSshPtyLiveness({ request, relayPtyId: PTY_ID })).resolves.toBe(expected)
|
||||
// One exact-id question, so the client neither scans an inventory nor interprets one.
|
||||
expect(request).toHaveBeenCalledTimes(1)
|
||||
})
|
||||
|
||||
it.each([
|
||||
['a relay too old to know the method', new Error('Method not found: pty.probeLiveness')],
|
||||
['a disposed multiplexer', new Error('Multiplexer disposed')],
|
||||
['a lost connection', new Error('SSH connection lost, reconnecting...')]
|
||||
])('fails closed on %s', async (_label, failure) => {
|
||||
const request = relay(failure)
|
||||
|
||||
await expect(probeSshPtyLiveness({ request, relayPtyId: PTY_ID })).resolves.toBeNull()
|
||||
})
|
||||
|
||||
it.each([
|
||||
['no answer at all', null],
|
||||
['an answer with no status', {}],
|
||||
['a status this client does not know', { status: 'probably_gone' }],
|
||||
['a non-string status', { status: true }]
|
||||
])('refuses to read %s as a verdict', async (_label, answer) => {
|
||||
const request = relay(answer)
|
||||
|
||||
await expect(probeSshPtyLiveness({ request, relayPtyId: PTY_ID })).resolves.toBeNull()
|
||||
})
|
||||
})
|
||||
@@ -0,0 +1,41 @@
|
||||
type RelayRequest = (
|
||||
method: string,
|
||||
params?: Record<string, unknown>,
|
||||
options?: { timeoutMs?: number }
|
||||
) => Promise<unknown>
|
||||
|
||||
const PROBE_TIMEOUT_MS = 10_000
|
||||
|
||||
/**
|
||||
* Forwards the liveness question to the relay, which is the only party that can answer it.
|
||||
*
|
||||
* The client deliberately computes nothing here. A previous version inferred an exit from an id's
|
||||
* absence in `pty.listProcesses` plus a matching mint epoch, and that is unsound: the relay also
|
||||
* removes a record without observing the process end, so a shutdown that gave up waiting for an
|
||||
* uninterruptible child produced a false death certificate (docs/reference/ssh-execution-boundary.md).
|
||||
*
|
||||
* Never throws, and fails closed. A relay too old to know the method answers JSON-RPC -32601, which
|
||||
* arrives here as a rejection and maps to null, exactly like a timeout or a disposed multiplexer.
|
||||
* Any status other than the two the owner certifies is unverifiable.
|
||||
*/
|
||||
export async function probeSshPtyLiveness(args: {
|
||||
request: RelayRequest
|
||||
relayPtyId: string
|
||||
}): Promise<boolean | null> {
|
||||
try {
|
||||
const answer = (await args.request(
|
||||
'pty.probeLiveness',
|
||||
{ id: args.relayPtyId },
|
||||
{ timeoutMs: PROBE_TIMEOUT_MS }
|
||||
)) as { status?: unknown } | null
|
||||
if (answer?.status === 'live') {
|
||||
return true
|
||||
}
|
||||
if (answer?.status === 'exited') {
|
||||
return false
|
||||
}
|
||||
return null
|
||||
} catch {
|
||||
return null
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,61 @@
|
||||
import { describe, expect, it, vi } from 'vitest'
|
||||
import { SshPtyProvider } from './ssh-pty-provider'
|
||||
import type { SshChannelMultiplexer } from '../ssh/ssh-channel-multiplexer'
|
||||
import { toAppSshPtyId } from '../../shared/ssh-pty-id'
|
||||
import { toRelayPtyIdWithMintEpoch } from '../../shared/relay-pty-mint-epoch'
|
||||
|
||||
const CONNECTION_ID = 'conn-1'
|
||||
const RELAY_EPOCH = '0f8f3a1e-1111-4111-8111-111111111111'
|
||||
const RELAY_PTY_ID = toRelayPtyIdWithMintEpoch(RELAY_EPOCH, 42)
|
||||
const APP_PTY_ID = toAppSshPtyId(CONNECTION_ID, RELAY_PTY_ID)
|
||||
|
||||
/**
|
||||
* The provider is the owner the liveness rule routes to, so this pins the contract seam itself:
|
||||
* without `probePtyLiveness` on this class, `probePtyLivenessFromRuntimeController` answers null
|
||||
* for every SSH id and no SSH worker can ever be certified exited. The verdict itself belongs to
|
||||
* the relay; this class only carries the question to it.
|
||||
*/
|
||||
function makeProvider(status: 'live' | 'exited' | 'unknown') {
|
||||
const request = vi.fn(async (method: string) => {
|
||||
if (method === 'pty.probeLiveness') {
|
||||
return { status }
|
||||
}
|
||||
throw new Error(`unexpected relay method ${method}`)
|
||||
})
|
||||
const mux = {
|
||||
request,
|
||||
onNotification: () => () => {},
|
||||
onRequest: () => () => {}
|
||||
} as unknown as SshChannelMultiplexer
|
||||
return { provider: new SshPtyProvider(CONNECTION_ID, mux), request }
|
||||
}
|
||||
|
||||
describe('SshPtyProvider liveness readback', () => {
|
||||
it('exposes the readback the liveness rule asks the owning provider for', () => {
|
||||
const { provider } = makeProvider('unknown')
|
||||
|
||||
expect(typeof provider.probePtyLiveness).toBe('function')
|
||||
})
|
||||
|
||||
it('answers live from the relay for an id its own cache has never seen', async () => {
|
||||
const { provider } = makeProvider('live')
|
||||
|
||||
expect(provider.hasPty(APP_PTY_ID)).toBe(false)
|
||||
await expect(provider.probePtyLiveness(APP_PTY_ID)).resolves.toBe(true)
|
||||
})
|
||||
|
||||
it('certifies an exit the relay observed', async () => {
|
||||
const { provider } = makeProvider('exited')
|
||||
|
||||
await expect(provider.probePtyLiveness(APP_PTY_ID)).resolves.toBe(false)
|
||||
})
|
||||
|
||||
it('refuses to answer for an id belonging to another SSH target', async () => {
|
||||
const { provider, request } = makeProvider('exited')
|
||||
|
||||
await expect(
|
||||
provider.probePtyLiveness(toAppSshPtyId('other-conn', RELAY_PTY_ID))
|
||||
).resolves.toBeNull()
|
||||
expect(request).not.toHaveBeenCalled()
|
||||
})
|
||||
})
|
||||
@@ -25,6 +25,7 @@ import { SshAgentSessionCapabilities } from './ssh-agent-session-capabilities'
|
||||
import type { PtyProcessInspection } from './pty-process-inspection'
|
||||
import { spawnWithTerminalRuntimeRepair, type TerminalRepairHook } from './ssh-pty-spawn-repair'
|
||||
import { createSshPtyProviderRpcOperations } from './ssh-pty-provider-rpc-operations'
|
||||
import { probeSshPtyLiveness } from './ssh-pty-liveness-probe'
|
||||
|
||||
// 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 {
|
||||
@@ -288,6 +289,21 @@ export class SshPtyProvider implements IPtyProvider {
|
||||
|
||||
hasPty = (id: string): boolean => this.livePtyIds.has(id)
|
||||
|
||||
/** `hasPty` is this process's cache of what it has seen; only the relay may answer absent. */
|
||||
probePtyLiveness = async (id: string): Promise<boolean | null> => {
|
||||
let relayPtyId: string
|
||||
try {
|
||||
relayPtyId = this.toRelayPtyId(id)
|
||||
} catch {
|
||||
// The id names another SSH target, so this provider is not its owner and cannot answer.
|
||||
return null
|
||||
}
|
||||
return await probeSshPtyLiveness({
|
||||
request: (method, params, options) => this.mux.request(method, params, options),
|
||||
relayPtyId
|
||||
})
|
||||
}
|
||||
|
||||
onData = (callback: SshPtyDataCallback): (() => void) => this.outputState.onData(callback)
|
||||
onRejectedData = (callback: SshPtyDataCallback): (() => void) =>
|
||||
this.outputState.onRejectedData(callback)
|
||||
|
||||
@@ -6,12 +6,15 @@ import type {
|
||||
LegacyWorkerRecoveryPorts,
|
||||
LegacyWorkerRecoveryResolution
|
||||
} from './runtime-legacy-worker-terminal-recovery-types'
|
||||
import { toAppSshPtyId } from '../../shared/ssh-pty-id'
|
||||
import { probeSshPtyLiveness } from '../providers/ssh-pty-liveness-probe'
|
||||
import { toAppSshPtyId, toRelaySshPtyId } from '../../shared/ssh-pty-id'
|
||||
import { toRelayPtyIdWithMintEpoch } from '../../shared/relay-pty-mint-epoch'
|
||||
|
||||
const CONNECTION_ID = 'conn-1'
|
||||
// A parseable SSH id, so the population this suite reasons about is the one it drives.
|
||||
const PTY_ID = toAppSshPtyId(CONNECTION_ID, 'pty-42')
|
||||
const OTHER_PTY_ID = toAppSshPtyId(CONNECTION_ID, 'pty-43')
|
||||
const RELAY_EPOCH = '0f8f3a1e-1111-4111-8111-111111111111'
|
||||
// A real SSH id, so the population this suite claims to cover is the one it drives.
|
||||
const PTY_ID = toAppSshPtyId(CONNECTION_ID, toRelayPtyIdWithMintEpoch(RELAY_EPOCH, 42))
|
||||
const OTHER_PTY_ID = toAppSshPtyId(CONNECTION_ID, toRelayPtyIdWithMintEpoch(RELAY_EPOCH, 43))
|
||||
|
||||
const candidate = {
|
||||
dispatchId: 'dispatch_1',
|
||||
@@ -48,7 +51,7 @@ const listingWithThePty = inventory({
|
||||
|
||||
function reconcile(options: {
|
||||
inventory: LegacyWorkerRecoveryInventory
|
||||
isPtyProvenAbsent: () => Promise<boolean>
|
||||
isPtyProvenAbsent: (ptyId: string) => Promise<boolean>
|
||||
refreshInventory?: LegacyWorkerRecoveryPorts['refreshInventory']
|
||||
}) {
|
||||
const pendingResolutions: LegacyWorkerRecoveryResolution[] = []
|
||||
@@ -149,14 +152,30 @@ describe('legacy worker recovery: certifying that a worker PTY exited', () => {
|
||||
expect(ports.isPtyProvenAbsent).not.toHaveBeenCalled()
|
||||
})
|
||||
|
||||
// No SSH provider implements the owner readback, so `isPtyProvenAbsent` is false for every SSH
|
||||
// candidate and the worker defers forever rather than being retired on a listing omission. That
|
||||
// is the honest verdict under the execution boundary and the deliberate cost of this change: the
|
||||
// operator clears such a worker with worker-stop, stop_unknown, then worker-abandon.
|
||||
it('defers an SSH candidate no owner can certify', async () => {
|
||||
// The rule `isLeafPtyProvenAbsent` applies: the owning provider's readback is proven absence
|
||||
// only when it answers false, and the SSH provider's readback is the relay itself.
|
||||
function provenAbsentViaRelay(request: Parameters<typeof probeSshPtyLiveness>[0]['request']) {
|
||||
return async (appPtyId: string): Promise<boolean> =>
|
||||
(await probeSshPtyLiveness({
|
||||
request,
|
||||
relayPtyId: toRelaySshPtyId(CONNECTION_ID, appPtyId)
|
||||
})) === false
|
||||
}
|
||||
|
||||
it('certifies an SSH exit the owning relay observed', async () => {
|
||||
const { pendingResolutions, deferredDispatchIds } = await reconcile({
|
||||
inventory: listingWithoutThePty,
|
||||
isPtyProvenAbsent: async () => false
|
||||
isPtyProvenAbsent: provenAbsentViaRelay(async () => ({ status: 'exited' }))
|
||||
})
|
||||
|
||||
expect(pendingResolutions).toEqual([{ candidate, resolution: 'exited' }])
|
||||
expect(deferredDispatchIds.size).toBe(0)
|
||||
})
|
||||
|
||||
it('defers an SSH candidate the owning relay will not certify', async () => {
|
||||
const { pendingResolutions, deferredDispatchIds } = await reconcile({
|
||||
inventory: listingWithoutThePty,
|
||||
isPtyProvenAbsent: provenAbsentViaRelay(async () => ({ status: 'unknown' }))
|
||||
})
|
||||
|
||||
expect(pendingResolutions).toEqual([])
|
||||
|
||||
@@ -0,0 +1,142 @@
|
||||
// Which removals of a PTY record observed the process end, and which are bookkeeping. Only the
|
||||
// relay knows, so only the relay may answer, and it certifies an exit from an observation alone.
|
||||
import { afterEach, beforeEach, describe, expect, it, vi } from 'vitest'
|
||||
|
||||
const { mockPtySpawn, mockPtyInstance, mockCreateShellPromptReadinessProbe } = vi.hoisted(() => ({
|
||||
mockPtySpawn: vi.fn(),
|
||||
mockCreateShellPromptReadinessProbe: vi.fn(),
|
||||
mockPtyInstance: {
|
||||
pid: process.pid,
|
||||
process: 'zsh',
|
||||
onData: vi.fn(),
|
||||
onExit: vi.fn(),
|
||||
write: vi.fn(),
|
||||
resize: vi.fn(),
|
||||
kill: vi.fn(),
|
||||
clear: vi.fn(),
|
||||
pause: vi.fn(),
|
||||
resume: vi.fn()
|
||||
}
|
||||
}))
|
||||
|
||||
vi.mock('node-pty', () => ({ spawn: mockPtySpawn }))
|
||||
|
||||
vi.mock('../main/pty/posix-pty-process-groups', () => ({
|
||||
forceKillPosixPtyProcessGroups: vi.fn((_pid: number, fallback: () => void) => fallback())
|
||||
}))
|
||||
|
||||
vi.mock('../main/shell-prompt-readiness-probe', () => ({
|
||||
createShellPromptReadinessProbe: mockCreateShellPromptReadinessProbe
|
||||
}))
|
||||
|
||||
import type { PtyHandler } from './pty-handler'
|
||||
import {
|
||||
beginPtyHandlerTest,
|
||||
createPtyRequestHelpers,
|
||||
endPtyHandlerTest,
|
||||
testPtyId
|
||||
} from './pty-handler-test-harness'
|
||||
import type { MockDispatcher } from './pty-handler-test-harness'
|
||||
|
||||
const UNREACHABLE_PID = 424_242
|
||||
|
||||
describe('PtyHandler.probeLiveness', () => {
|
||||
let dispatcher: MockDispatcher
|
||||
let handler: PtyHandler
|
||||
let originalPlatform: PropertyDescriptor | undefined
|
||||
|
||||
const { spawnPty } = createPtyRequestHelpers(() => dispatcher)
|
||||
|
||||
beforeEach(() => {
|
||||
;({ dispatcher, handler, originalPlatform } = beginPtyHandlerTest({
|
||||
mockPtySpawn,
|
||||
mockPtyInstance,
|
||||
mockCreateShellPromptReadinessProbe
|
||||
}))
|
||||
})
|
||||
|
||||
afterEach(async () => {
|
||||
await endPtyHandlerTest(handler, originalPlatform)
|
||||
})
|
||||
|
||||
async function probe(id: string): Promise<string> {
|
||||
const answer = (await dispatcher.callRequest('pty.probeLiveness', { id })) as {
|
||||
status: string
|
||||
}
|
||||
return answer.status
|
||||
}
|
||||
|
||||
describe('observed exits', () => {
|
||||
it('answers live while the relay owns a running record', async () => {
|
||||
const { id } = await spawnPty()
|
||||
|
||||
expect(await probe(id)).toBe('live')
|
||||
})
|
||||
|
||||
it('certifies the exit node-pty reported', async () => {
|
||||
const { id } = await spawnPty()
|
||||
const onExit = mockPtyInstance.onExit.mock.calls.at(-1)?.[0] as (event: {
|
||||
exitCode: number
|
||||
signal?: number
|
||||
}) => void
|
||||
onExit({ exitCode: 0 })
|
||||
|
||||
expect(await probe(id)).toBe('exited')
|
||||
})
|
||||
|
||||
it('certifies an exit it proves by probing the pid, and keeps certifying it', async () => {
|
||||
mockPtySpawn.mockReturnValueOnce({ ...mockPtyInstance, pid: UNREACHABLE_PID })
|
||||
const { id } = await spawnPty()
|
||||
|
||||
// The shell ended without node-pty reporting it; the pid answers ESRCH.
|
||||
expect(await probe(id)).toBe('exited')
|
||||
// The first probe reaped the record, so the second has only the observation to go on — and
|
||||
// the recovery sweep that drives this asks again on every pass.
|
||||
expect(await probe(id)).toBe('exited')
|
||||
})
|
||||
})
|
||||
|
||||
describe('bookkeeping removals and states that observed nothing', () => {
|
||||
it('stays unverifiable after a shutdown removes a record it never saw exit', async () => {
|
||||
const { id } = await spawnPty()
|
||||
const disposal = handler.dispose()
|
||||
await vi.advanceTimersByTimeAsync(8_001)
|
||||
await disposal
|
||||
|
||||
expect(handler.activePtyCount).toBe(0)
|
||||
expect(await probe(id)).toBe('unknown')
|
||||
})
|
||||
|
||||
it('stays unverifiable for an id this relay never minted', async () => {
|
||||
expect(await probe(testPtyId(99))).toBe('unknown')
|
||||
})
|
||||
|
||||
it('stays unverifiable for an id another relay generation minted', async () => {
|
||||
expect(await probe('pty2:some-other-generation:1')).toBe('unknown')
|
||||
})
|
||||
|
||||
it.each(['pty-7', '', 'not-a-pty-id'])(
|
||||
'stays unverifiable for the id shape %s, which names no generation',
|
||||
async (id) => {
|
||||
expect(await probe(id)).toBe('unknown')
|
||||
}
|
||||
)
|
||||
|
||||
it('separates a real exit from a teardown removal in the same shutdown', async () => {
|
||||
const exiting = await spawnPty()
|
||||
const onExit = mockPtyInstance.onExit.mock.calls.at(-1)?.[0] as (event: {
|
||||
exitCode: number
|
||||
}) => void
|
||||
onExit({ exitCode: 0 })
|
||||
const tornDown = await spawnPty()
|
||||
|
||||
const disposal = handler.dispose()
|
||||
await vi.advanceTimersByTimeAsync(8_001)
|
||||
await disposal
|
||||
|
||||
// Both records are gone from the map; only one of them was ever watched ending.
|
||||
expect(await probe(exiting.id)).toBe('exited')
|
||||
expect(await probe(tornDown.id)).toBe('unknown')
|
||||
})
|
||||
})
|
||||
})
|
||||
@@ -0,0 +1,141 @@
|
||||
// A shutdown timeout removes bookkeeping without proving the process exited.
|
||||
import { describe, expect, it, vi, beforeEach, afterEach } from 'vitest'
|
||||
|
||||
const { mockPtySpawn, mockPtyInstance, mockCreateShellPromptReadinessProbe } = vi.hoisted(() => ({
|
||||
mockPtySpawn: vi.fn(),
|
||||
mockCreateShellPromptReadinessProbe: vi.fn(),
|
||||
mockPtyInstance: {
|
||||
pid: process.pid,
|
||||
process: 'zsh',
|
||||
onData: vi.fn(),
|
||||
onExit: vi.fn(),
|
||||
write: vi.fn(),
|
||||
resize: vi.fn(),
|
||||
kill: vi.fn(),
|
||||
clear: vi.fn(),
|
||||
pause: vi.fn(),
|
||||
resume: vi.fn()
|
||||
}
|
||||
}))
|
||||
|
||||
vi.mock('node-pty', () => ({ spawn: mockPtySpawn }))
|
||||
|
||||
vi.mock('../main/pty/posix-pty-process-groups', () => ({
|
||||
forceKillPosixPtyProcessGroups: vi.fn((_pid: number, fallback: () => void) => fallback())
|
||||
}))
|
||||
|
||||
vi.mock('../main/shell-prompt-readiness-probe', () => ({
|
||||
createShellPromptReadinessProbe: mockCreateShellPromptReadinessProbe
|
||||
}))
|
||||
|
||||
import type { PtyHandler } from './pty-handler'
|
||||
import { SshPtyProvider } from '../main/providers/ssh-pty-provider'
|
||||
import { toAppSshPtyId } from '../shared/ssh-pty-id'
|
||||
import { reconcileLegacyWorkerCandidate } from '../main/runtime/runtime-legacy-worker-terminal-recovery-candidate'
|
||||
import type { LegacyWorkerRecoveryResolution } from '../main/runtime/runtime-legacy-worker-terminal-recovery-types'
|
||||
import {
|
||||
beginPtyHandlerTest,
|
||||
createPtyRequestHelpers,
|
||||
endPtyHandlerTest
|
||||
} from './pty-handler-test-harness'
|
||||
import type { MockDispatcher } from './pty-handler-test-harness'
|
||||
import { parseRelayPtyMintEpoch } from '../shared/relay-pty-mint-epoch'
|
||||
|
||||
describe('SSH exit certification across relay shutdown', () => {
|
||||
let dispatcher: MockDispatcher
|
||||
let handler: PtyHandler
|
||||
let originalPlatform: PropertyDescriptor | undefined
|
||||
|
||||
const { spawnPty } = createPtyRequestHelpers(() => dispatcher)
|
||||
|
||||
beforeEach(() => {
|
||||
;({ dispatcher, handler, originalPlatform } = beginPtyHandlerTest({
|
||||
mockPtySpawn,
|
||||
mockPtyInstance,
|
||||
mockCreateShellPromptReadinessProbe
|
||||
}))
|
||||
})
|
||||
|
||||
afterEach(async () => {
|
||||
await endPtyHandlerTest(handler, originalPlatform)
|
||||
})
|
||||
|
||||
async function capabilities(): Promise<{ ptyIdMintEpoch?: unknown }> {
|
||||
return (await dispatcher.callRequest('pty.getCapabilities', {})) as { ptyIdMintEpoch?: unknown }
|
||||
}
|
||||
|
||||
it('names the generation that minted its PTY ids', async () => {
|
||||
const { id } = await spawnPty()
|
||||
|
||||
const mintEpoch = parseRelayPtyMintEpoch(id)
|
||||
expect(mintEpoch).toBeTruthy()
|
||||
expect((await capabilities()).ptyIdMintEpoch).toBe(mintEpoch)
|
||||
})
|
||||
|
||||
it('keeps that generation stable across ids and reads', async () => {
|
||||
const first = await spawnPty()
|
||||
const second = await spawnPty()
|
||||
|
||||
expect(parseRelayPtyMintEpoch(second.id)).toBe(parseRelayPtyMintEpoch(first.id))
|
||||
expect((await capabilities()).ptyIdMintEpoch).toBe((await capabilities()).ptyIdMintEpoch)
|
||||
})
|
||||
it('does not certify exit after shutdown bookkeeping removes an unexited PTY', async () => {
|
||||
const { id } = await spawnPty()
|
||||
const provider = new SshPtyProvider('review-target', {
|
||||
request: (method: string, params: Record<string, unknown>) =>
|
||||
dispatcher.callRequest(method, params),
|
||||
onNotification: () => () => {},
|
||||
onRequest: () => () => {}
|
||||
} as never)
|
||||
const appId = toAppSshPtyId('review-target', id)
|
||||
expect(await provider.probePtyLiveness(appId)).toBe(true)
|
||||
const disposal = handler.dispose()
|
||||
await vi.advanceTimersByTimeAsync(8_001)
|
||||
await disposal
|
||||
expect(dispatcher._notifications.filter((event) => event.method === 'pty.exit')).toEqual([])
|
||||
expect(mockPtyInstance.kill).toHaveBeenCalled()
|
||||
expect(() => process.kill(process.pid, 0)).not.toThrow()
|
||||
expect(handler.activePtyCount).toBe(0)
|
||||
expect(await provider.probePtyLiveness(appId)).toBeNull()
|
||||
})
|
||||
|
||||
it('keeps missing shutdown records unverifiable when another failed kill leaves the relay serving', async () => {
|
||||
const { id } = await spawnPty()
|
||||
const failedKill = vi.fn<() => void>(() => {
|
||||
throw new Error('host refused kill')
|
||||
})
|
||||
mockPtySpawn.mockReturnValueOnce({ ...mockPtyInstance, kill: failedKill })
|
||||
await spawnPty()
|
||||
const disposal = handler.dispose().catch((error: Error) => error)
|
||||
await vi.advanceTimersByTimeAsync(8_001)
|
||||
expect(await disposal).toMatchObject({ message: 'host refused kill' })
|
||||
expect(handler.activePtyCount).toBe(1)
|
||||
const provider = new SshPtyProvider('review-target', {
|
||||
request: (method: string, params: Record<string, unknown>) =>
|
||||
dispatcher.callRequest(method, params),
|
||||
onNotification: () => () => {},
|
||||
onRequest: () => () => {}
|
||||
} as never)
|
||||
expect(() => process.kill(process.pid, 0)).not.toThrow()
|
||||
expect(dispatcher._notifications.filter((event) => event.method === 'pty.exit')).toEqual([])
|
||||
// Permit cleanup of the second fixture after the assertion.
|
||||
failedKill.mockImplementation(() => {})
|
||||
const appId = toAppSshPtyId('review-target', id)
|
||||
const verdict = await provider.probePtyLiveness(appId)
|
||||
const pendingResolutions: LegacyWorkerRecoveryResolution[] = []
|
||||
const deferredDispatchIds = new Set<string>()
|
||||
await reconcileLegacyWorkerCandidate({
|
||||
controller: {} as never,
|
||||
ports: { isPtyProvenAbsent: async () => verdict === false } as never,
|
||||
options: {},
|
||||
candidate: { dispatchId: 'review-worker', ptyId: appId } as never,
|
||||
workspace: {} as never,
|
||||
resolvedWorktrees: [],
|
||||
inventory: { livePtyIds: new Set() } as never,
|
||||
pendingResolutions,
|
||||
deferredDispatchIds
|
||||
})
|
||||
expect(pendingResolutions).toEqual([])
|
||||
expect([...deferredDispatchIds]).toEqual(['review-worker'])
|
||||
})
|
||||
})
|
||||
@@ -36,6 +36,7 @@ import {
|
||||
} from './pty-spawn-cwd'
|
||||
import { PhysicalExitTracker } from '../shared/physical-exit-tracker'
|
||||
import { PTY_ATTACH_PROVEN_EXITED_MARKER } from '../shared/pty-attach-absence-evidence'
|
||||
import { toRelayPtyIdWithMintEpoch } from '../shared/relay-pty-mint-epoch'
|
||||
import { SHELL_READY_MARKER_PREFIX } from '../main/shell-ready-marker-scanner'
|
||||
import {
|
||||
createShellStartupOutputScanState,
|
||||
@@ -353,6 +354,9 @@ const STARTUP_COMMAND_SHELL_READY_FALLBACK_MS = 1500
|
||||
const RENDERER_SHELL_READY_RETENTION_MS = 15_000
|
||||
const PTY_FORCE_KILL_RETRY_DELAY_MS = 250
|
||||
const PTY_FORCE_KILL_MAX_ATTEMPTS = 2
|
||||
|
||||
/** Cap on remembered observed exits. Eviction answers unverifiable, never a false exit. */
|
||||
const OBSERVED_PTY_EXIT_HISTORY = 4_096
|
||||
const ALLOWED_SIGNALS = new Set([
|
||||
'SIGINT',
|
||||
'SIGTERM',
|
||||
@@ -533,6 +537,15 @@ export class PtyHandler {
|
||||
private interactiveOutputCharsByPty = new Map<string, number>()
|
||||
private pendingSpawnCount = 0
|
||||
private pendingReviveIds = new Set<string>()
|
||||
/**
|
||||
* Ids this relay retired having OBSERVED the process end: node-pty's `onExit`, or a pid probe
|
||||
* that answered ESRCH. Every other removal is bookkeeping — a shutdown that gave up waiting for
|
||||
* an uninterruptible child still deletes the record — so absence from `this.ptys` cannot be read
|
||||
* as an exit. This is the positive half, recorded where the observation happens, and a liveness
|
||||
* probe certifies an exit from nothing else. Bounded: an evicted id answers unverifiable, which
|
||||
* is the safe direction (docs/reference/ssh-execution-boundary.md).
|
||||
*/
|
||||
private readonly observedPtyExitIds = new Set<string>()
|
||||
private creationFenced = false
|
||||
private pendingCreationDrainResolvers = new Set<() => void>()
|
||||
private worktreeRemovalCoordinator: RelayPtyWorktreeRemovalCoordinator | null = null
|
||||
@@ -765,6 +778,17 @@ export class PtyHandler {
|
||||
}
|
||||
}
|
||||
|
||||
/** Called only where the process end was observed; see {@link observedPtyExitIds}. */
|
||||
private recordObservedPtyExit(id: string): void {
|
||||
if (this.observedPtyExitIds.size >= OBSERVED_PTY_EXIT_HISTORY) {
|
||||
const oldest = this.observedPtyExitIds.values().next()
|
||||
if (!oldest.done) {
|
||||
this.observedPtyExitIds.delete(oldest.value)
|
||||
}
|
||||
}
|
||||
this.observedPtyExitIds.add(id)
|
||||
}
|
||||
|
||||
// Why: the sole removal path, so the three exit routes can't drift on who announces an empty pool.
|
||||
private removePty(id: string): void {
|
||||
this.ptys.delete(id)
|
||||
@@ -1039,6 +1063,7 @@ export class PtyHandler {
|
||||
this.publishPendingExit(managed.id)
|
||||
this.notifyExitListener(managed)
|
||||
this.agentSessionOwners.release(managed.id)
|
||||
this.recordObservedPtyExit(managed.id)
|
||||
this.removePty(managed.id)
|
||||
this.clearPtyInputState(managed.id)
|
||||
// Why: release the ptmx fd on natural exit, else the master fd leaks until GC (docs/fix-pty-fd-leak.md).
|
||||
@@ -1092,8 +1117,13 @@ export class PtyHandler {
|
||||
agentSessionCreateOperationVersion: AGENT_SESSION_CREATE_OPERATION_PROTOCOL_VERSION,
|
||||
// Additive capability: clients may request the no-process-table inventory
|
||||
// projection and consume fenced inspect evidence on this host.
|
||||
foregroundProcessEvidenceVersion: 1
|
||||
foregroundProcessEvidenceVersion: 1,
|
||||
// Additive: names the generation that minted this relay's ids, so a client can read an id's
|
||||
// absence from `pty.listProcesses` as an exit this host observed rather than as a restart.
|
||||
// A relay that omits it leaves every absence unverifiable, which is the shipped behaviour.
|
||||
ptyIdMintEpoch: this.ptyIdMintEpoch
|
||||
}))
|
||||
this.dispatcher.onRequest('pty.probeLiveness', async (params) => this.probeLiveness(params))
|
||||
this.dispatcher.onRequest('pty.listProcesses', (params) => this.listProcesses(params))
|
||||
this.dispatcher.onRequest('pty.getDefaultShell', async () => resolveDefaultShell())
|
||||
this.dispatcher.onRequest('pty.serialize', (p) => this.serialize(p))
|
||||
@@ -1863,7 +1893,7 @@ export class PtyHandler {
|
||||
const shell = resolvedShellOverride || requestedEnvShell || resolveDefaultShell()
|
||||
let id: string
|
||||
do {
|
||||
id = `pty2:${encodeURIComponent(this.ptyIdMintEpoch)}:${this.nextId++}`
|
||||
id = toRelayPtyIdWithMintEpoch(this.ptyIdMintEpoch, this.nextId++)
|
||||
} while (this.ptys.has(id) || this.pendingReviveIds.has(id))
|
||||
|
||||
// Why: augmenter values override renderer env so remote paths and hook coords win over local userData.
|
||||
@@ -2464,6 +2494,9 @@ export class PtyHandler {
|
||||
if (!managed.pty.pid || isProcessAlive(managed.pty.pid)) {
|
||||
return false
|
||||
}
|
||||
// Recorded here rather than inside the shared reap: this is the line that observed the pid was
|
||||
// gone. The reap's other caller is tidying a torn-down record and watched nothing.
|
||||
this.recordObservedPtyExit(managed.id)
|
||||
this.reapExitedPty(managed, 'exited')
|
||||
return true
|
||||
}
|
||||
@@ -2600,6 +2633,38 @@ export class PtyHandler {
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* Exact-id liveness, answered from this relay's own lifecycle state. The client cannot compute
|
||||
* this: absence from an inventory means "no record", and only the owner knows whether a record
|
||||
* left because the process ended or because teardown stopped waiting for it
|
||||
* (docs/reference/ssh-execution-boundary.md).
|
||||
*
|
||||
* `exited` requires an observation — a live record whose pid probes absent, or an id in
|
||||
* {@link observedPtyExitIds}. Everything else is `unknown`: an id this generation never minted
|
||||
* or never saw end, a revive in flight, a record mid-teardown, and any id at all once a shutdown
|
||||
* has been fenced, because that is when records start leaving without proof.
|
||||
*/
|
||||
private async probeLiveness(
|
||||
params: Record<string, unknown>
|
||||
): Promise<{ status: 'live' | 'exited' | 'unknown' }> {
|
||||
const id = typeof params.id === 'string' ? params.id : ''
|
||||
const managed = this.ptys.get(id)
|
||||
if (managed) {
|
||||
if (managed.disposed) {
|
||||
return { status: 'unknown' }
|
||||
}
|
||||
return { status: this.reapPtyProvenExited(managed) ? 'exited' : 'live' }
|
||||
}
|
||||
if (this.pendingReviveIds.has(id)) {
|
||||
return { status: 'unknown' }
|
||||
}
|
||||
// The ledger is the only gate, deliberately. A blanket "shutdown has started" refusal would
|
||||
// also mask a wrong entry in it, since teardown only ever runs behind that fence — and it
|
||||
// would withhold exits this relay really did observe. Keyed by the whole id, which embeds this
|
||||
// generation, so an id another generation minted can never be in it.
|
||||
return { status: this.observedPtyExitIds.has(id) ? 'exited' : 'unknown' }
|
||||
}
|
||||
|
||||
private async hasChildProcesses(params: Record<string, unknown>): Promise<boolean> {
|
||||
const id = params.id as string
|
||||
const managed = this.ptys.get(id)
|
||||
|
||||
@@ -0,0 +1,30 @@
|
||||
/**
|
||||
* A relay PTY id carries the mint epoch of the relay process that allocated it, so a client can
|
||||
* tell the two halves of "this relay does not list that id" apart: the relay minted it and no
|
||||
* longer has it, which is an exit the host observed, versus the relay never had it, which is every
|
||||
* id minted before a restart and is evidence of nothing (docs/reference/ssh-execution-boundary.md).
|
||||
*
|
||||
* The id shape lives here so the relay that mints and the client that reads it cannot drift.
|
||||
*/
|
||||
const MINT_EPOCH_PTY_ID_PREFIX = 'pty2:'
|
||||
|
||||
export function toRelayPtyIdWithMintEpoch(mintEpoch: string, sequence: number): string {
|
||||
return `${MINT_EPOCH_PTY_ID_PREFIX}${encodeURIComponent(mintEpoch)}:${sequence}`
|
||||
}
|
||||
|
||||
/** Null for a legacy `pty-N` id, which names no epoch and so can never certify an exit. */
|
||||
export function parseRelayPtyMintEpoch(relayPtyId: string): string | null {
|
||||
if (!relayPtyId.startsWith(MINT_EPOCH_PTY_ID_PREFIX)) {
|
||||
return null
|
||||
}
|
||||
const remainder = relayPtyId.slice(MINT_EPOCH_PTY_ID_PREFIX.length)
|
||||
const sequenceSeparator = remainder.lastIndexOf(':')
|
||||
if (sequenceSeparator <= 0) {
|
||||
return null
|
||||
}
|
||||
try {
|
||||
return decodeURIComponent(remainder.slice(0, sequenceSeparator)) || null
|
||||
} catch {
|
||||
return null
|
||||
}
|
||||
}
|
||||
Reference in New Issue
Block a user