mirror of
https://github.com/stablyai/orca.git
synced 2026-09-30 00:03:15 +00:00
fix(relay): let the owner say why a PTY record disappeared
The SSH liveness probe read "absent from pty.listProcesses and minted by this relay's generation" as an observed exit. That implication is false. Relay shutdown waits eight seconds for a physical exit, catches the timeout, and deletes the record anyway; a sibling PTY's failed kill rejects the aggregate, the grace lifecycle defers shutdown, and the relay keeps serving under the same mint epoch. A provider that never attached the PTY then received a death certificate for a running process, and recovery queued its worker as exited. The census of every record removal splits cleanly: node-pty's onExit and the ESRCH pid probe watched the process end; the shutdown teardown and the torn-down-record sweep did not. Only the first two now record the id, at the line that made the observation rather than in the shared reap they both pass through. pty.probeLiveness answers an exact id from that state: a live record, or exited when its pid probes absent or the id is in the ledger. A record mid-teardown, a revive in flight, and every id the relay never watched end are unknown. It fails closed, so a removal path added later without a record defers instead of certifying. SshPtyProvider forwards and maps three statuses; an older relay answers -32601, which arrives as a rejection and maps to null, so no capability negotiation is needed. Two client-side gates went with it. The mint-epoch comparison is subsumed by ledger keys that already embed the generation, and a blanket shutdown-fenced refusal only masked the ledger in the one path where teardown runs.
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.
|
||||
@@ -1,104 +1,50 @@
|
||||
import { describe, expect, it, vi } from 'vitest'
|
||||
import { probeSshPtyLiveness } from './ssh-pty-liveness-probe'
|
||||
import {
|
||||
parseRelayPtyMintEpoch,
|
||||
toRelayPtyIdWithMintEpoch
|
||||
} from '../../shared/relay-pty-mint-epoch'
|
||||
|
||||
const EPOCH = '0f8f3a1e-1111-4111-8111-111111111111'
|
||||
const PTY_ID = toRelayPtyIdWithMintEpoch(EPOCH, 7)
|
||||
const PTY_ID = 'pty2:0f8f3a1e-1111-4111-8111-111111111111:7'
|
||||
|
||||
function relay(options: {
|
||||
listed?: { id: string }[] | null
|
||||
epoch?: string | null
|
||||
listThrows?: Error
|
||||
capabilitiesThrow?: Error
|
||||
}) {
|
||||
return vi.fn(async (method: string) => {
|
||||
if (method === 'pty.listProcesses') {
|
||||
if (options.listThrows) {
|
||||
throw options.listThrows
|
||||
}
|
||||
return options.listed === undefined ? [] : options.listed
|
||||
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
|
||||
}
|
||||
if (method === 'pty.getCapabilities') {
|
||||
if (options.capabilitiesThrow) {
|
||||
throw options.capabilitiesThrow
|
||||
}
|
||||
return options.epoch === null ? {} : { ptyIdMintEpoch: options.epoch ?? EPOCH }
|
||||
}
|
||||
throw new Error(`unexpected relay method ${method}`)
|
||||
return answer
|
||||
})
|
||||
}
|
||||
|
||||
describe('SSH PTY liveness probe', () => {
|
||||
it('answers live for an id the relay still lists', async () => {
|
||||
const request = relay({ listed: [{ id: PTY_ID }] })
|
||||
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(true)
|
||||
// A listed id needs no epoch comparison, so the second round trip is not spent.
|
||||
expect(request).toHaveBeenCalledTimes(1)
|
||||
})
|
||||
|
||||
it('certifies the exit when the relay that minted the id no longer lists it', async () => {
|
||||
const request = relay({ listed: [{ id: toRelayPtyIdWithMintEpoch(EPOCH, 8) }] })
|
||||
|
||||
await expect(probeSshPtyLiveness({ request, relayPtyId: PTY_ID })).resolves.toBe(false)
|
||||
})
|
||||
|
||||
it('certifies the exit even when the relay now lists nothing at all', async () => {
|
||||
// The population the recovery sweep meets: the worker was the host's only terminal and its
|
||||
// shell died while Orca was closed, so no listed id can name the current generation.
|
||||
const request = relay({ listed: [] })
|
||||
|
||||
await expect(probeSshPtyLiveness({ request, relayPtyId: PTY_ID })).resolves.toBe(false)
|
||||
})
|
||||
|
||||
it('stays unverifiable when a restarted relay disowns an id it never minted', async () => {
|
||||
const request = relay({ listed: [], epoch: 'ffffffff-2222-4222-8222-222222222222' })
|
||||
|
||||
await expect(probeSshPtyLiveness({ request, relayPtyId: PTY_ID })).resolves.toBeNull()
|
||||
})
|
||||
|
||||
it('stays unverifiable for a relay that names no mint epoch', async () => {
|
||||
const request = relay({ listed: [], epoch: null })
|
||||
|
||||
await expect(probeSshPtyLiveness({ request, relayPtyId: PTY_ID })).resolves.toBeNull()
|
||||
})
|
||||
|
||||
it('stays unverifiable for a legacy id that carries no epoch', async () => {
|
||||
const request = relay({ listed: [] })
|
||||
|
||||
await expect(probeSshPtyLiveness({ request, relayPtyId: 'pty-4' })).resolves.toBeNull()
|
||||
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([
|
||||
['the listing', { listThrows: new Error('Multiplexer disposed') }],
|
||||
['the capability read', { capabilitiesThrow: new Error('SSH connection lost') }]
|
||||
])('stays unverifiable when %s fails', async (_label, failure) => {
|
||||
const request = relay({ listed: [], ...failure })
|
||||
['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('stays unverifiable when the relay answers with no listing at all', async () => {
|
||||
const request = relay({ listed: null })
|
||||
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()
|
||||
})
|
||||
})
|
||||
|
||||
describe('relay PTY mint epoch', () => {
|
||||
it('round-trips an epoch through an id', () => {
|
||||
expect(parseRelayPtyMintEpoch(toRelayPtyIdWithMintEpoch(EPOCH, 3))).toBe(EPOCH)
|
||||
})
|
||||
|
||||
it('round-trips an epoch that needs escaping', () => {
|
||||
expect(parseRelayPtyMintEpoch(toRelayPtyIdWithMintEpoch('a:b/c d', 1))).toBe('a:b/c d')
|
||||
})
|
||||
|
||||
it.each(['pty-4', 'pty2:', 'pty2::9', 'shell-1'])('names no epoch in %s', (id) => {
|
||||
expect(parseRelayPtyMintEpoch(id)).toBeNull()
|
||||
})
|
||||
})
|
||||
|
||||
@@ -1,5 +1,3 @@
|
||||
import { parseRelayPtyMintEpoch } from '../../shared/relay-pty-mint-epoch'
|
||||
|
||||
type RelayRequest = (
|
||||
method: string,
|
||||
params?: Record<string, unknown>,
|
||||
@@ -9,51 +7,34 @@ type RelayRequest = (
|
||||
const PROBE_TIMEOUT_MS = 10_000
|
||||
|
||||
/**
|
||||
* Whether the relay still owns this PTY, as an answer a caller may retire durable state on.
|
||||
* Forwards the liveness question to the relay, which is the only party that can answer it.
|
||||
*
|
||||
* `pty.listProcesses` is the relay's own live set: it reaps a torn-down record and re-probes a pid
|
||||
* before reporting, so a listed id is live. An unlisted id is only an exit when this relay is also
|
||||
* the one that minted it — after a restart the relay disowns every id the previous one allocated
|
||||
* without having checked anything, which is why absence alone was never evidence
|
||||
* (docs/reference/ssh-execution-boundary.md).
|
||||
* 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).
|
||||
*
|
||||
* Deliberately not `pty.attach`, the only refusal that carries the proven-exited marker: a
|
||||
* successful attach opens a delivery, retires the previous one, and restores retired pane surfaces,
|
||||
* so probing with it would disturb a live consumer in exactly the case where the answer is "live".
|
||||
*
|
||||
* Never throws. A transport failure, a disposed multiplexer, a timeout, a legacy `pty-N` id, and a
|
||||
* relay that names no mint epoch all answer null, because none of them observed the process.
|
||||
* 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 listed = (await args.request(
|
||||
'pty.listProcesses',
|
||||
{ includeForegroundProcessEvidence: false },
|
||||
const answer = (await args.request(
|
||||
'pty.probeLiveness',
|
||||
{ id: args.relayPtyId },
|
||||
{ timeoutMs: PROBE_TIMEOUT_MS }
|
||||
)) as { id?: unknown }[] | null
|
||||
if (!Array.isArray(listed)) {
|
||||
return null
|
||||
}
|
||||
if (listed.some((session) => session.id === args.relayPtyId)) {
|
||||
)) as { status?: unknown } | null
|
||||
if (answer?.status === 'live') {
|
||||
return true
|
||||
}
|
||||
const mintEpoch = parseRelayPtyMintEpoch(args.relayPtyId)
|
||||
if (!mintEpoch) {
|
||||
return null
|
||||
if (answer?.status === 'exited') {
|
||||
return false
|
||||
}
|
||||
// Read after the listing on purpose: both answers come from one relay process, and a restart
|
||||
// between them breaks the multiplexer rather than pairing a new epoch with an old listing.
|
||||
const capabilities = (await args.request('pty.getCapabilities', undefined, {
|
||||
timeoutMs: PROBE_TIMEOUT_MS
|
||||
})) as { ptyIdMintEpoch?: unknown } | null
|
||||
const currentEpoch = capabilities?.ptyIdMintEpoch
|
||||
if (typeof currentEpoch !== 'string' || currentEpoch !== mintEpoch) {
|
||||
return null
|
||||
}
|
||||
return false
|
||||
return null
|
||||
} catch {
|
||||
return null
|
||||
}
|
||||
|
||||
@@ -12,15 +12,13 @@ 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.
|
||||
* 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(listed: { id: string }[]) {
|
||||
function makeProvider(status: 'live' | 'exited' | 'unknown') {
|
||||
const request = vi.fn(async (method: string) => {
|
||||
if (method === 'pty.listProcesses') {
|
||||
return listed
|
||||
}
|
||||
if (method === 'pty.getCapabilities') {
|
||||
return { ptyIdMintEpoch: RELAY_EPOCH }
|
||||
if (method === 'pty.probeLiveness') {
|
||||
return { status }
|
||||
}
|
||||
throw new Error(`unexpected relay method ${method}`)
|
||||
})
|
||||
@@ -34,26 +32,26 @@ function makeProvider(listed: { id: string }[]) {
|
||||
|
||||
describe('SshPtyProvider liveness readback', () => {
|
||||
it('exposes the readback the liveness rule asks the owning provider for', () => {
|
||||
const { provider } = makeProvider([])
|
||||
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([{ id: RELAY_PTY_ID }])
|
||||
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([])
|
||||
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([])
|
||||
const { provider, request } = makeProvider('exited')
|
||||
|
||||
await expect(
|
||||
provider.probePtyLiveness(toAppSshPtyId('other-conn', RELAY_PTY_ID))
|
||||
|
||||
@@ -162,24 +162,20 @@ describe('legacy worker recovery: certifying that a worker PTY exited', () => {
|
||||
})) === false
|
||||
}
|
||||
|
||||
it('certifies an SSH exit when the relay that minted the id no longer lists it', async () => {
|
||||
it('certifies an SSH exit the owning relay observed', async () => {
|
||||
const { pendingResolutions, deferredDispatchIds } = await reconcile({
|
||||
inventory: listingWithoutThePty,
|
||||
isPtyProvenAbsent: provenAbsentViaRelay(async (method) =>
|
||||
method === 'pty.listProcesses' ? [] : { ptyIdMintEpoch: RELAY_EPOCH }
|
||||
)
|
||||
isPtyProvenAbsent: provenAbsentViaRelay(async () => ({ status: 'exited' }))
|
||||
})
|
||||
|
||||
expect(pendingResolutions).toEqual([{ candidate, resolution: 'exited' }])
|
||||
expect(deferredDispatchIds.size).toBe(0)
|
||||
})
|
||||
|
||||
it('defers an SSH candidate a restarted relay merely disowns', async () => {
|
||||
it('defers an SSH candidate the owning relay will not certify', async () => {
|
||||
const { pendingResolutions, deferredDispatchIds } = await reconcile({
|
||||
inventory: listingWithoutThePty,
|
||||
isPtyProvenAbsent: provenAbsentViaRelay(async (method) =>
|
||||
method === 'pty.listProcesses' ? [] : { ptyIdMintEpoch: 'a-later-relay-generation' }
|
||||
)
|
||||
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')
|
||||
})
|
||||
})
|
||||
})
|
||||
@@ -1,82 +0,0 @@
|
||||
// The relay half of the SSH liveness readback. A client may only read "this id is absent from my
|
||||
// listing" as an exit when the relay that minted the id is the one answering, so the generation
|
||||
// stamped into every id has to be the generation `pty.getCapabilities` names. If these two ever
|
||||
// disagree, every SSH absence silently becomes unverifiable forever and no SSH worker can be
|
||||
// retired (docs/reference/ssh-execution-boundary.md).
|
||||
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 {
|
||||
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('PtyHandler mint epoch', () => {
|
||||
let dispatcher: MockDispatcher
|
||||
let handler: PtyHandler
|
||||
let originalPlatform: PropertyDescriptor | undefined
|
||||
|
||||
const { spawnPty } = createPtyRequestHelpers(() => dispatcher)
|
||||
|
||||
beforeEach(() => {
|
||||
;({ dispatcher, handler, originalPlatform } = beginPtyHandlerTest({
|
||||
mockPtySpawn,
|
||||
mockPtyInstance,
|
||||
mockCreateShellPromptReadinessProbe
|
||||
}))
|
||||
})
|
||||
|
||||
afterEach(() => {
|
||||
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)
|
||||
})
|
||||
})
|
||||
@@ -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'])
|
||||
})
|
||||
})
|
||||
@@ -354,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',
|
||||
@@ -534,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
|
||||
@@ -766,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)
|
||||
@@ -1040,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).
|
||||
@@ -1099,6 +1123,7 @@ export class PtyHandler {
|
||||
// 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))
|
||||
@@ -2469,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
|
||||
}
|
||||
@@ -2605,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)
|
||||
|
||||
Reference in New Issue
Block a user