mirror of
https://github.com/stablyai/orca.git
synced 2026-09-29 00:02:56 +00:00
fix(daemon): hold an observed session exit until its owner receives it
The daemon's only death evidence was `retiredIncarnations`, a 2-second tombstone written in `onSessionExit` from a session the host reaps in the same callback. Two seconds is a window no reopening app can hit, so once main stopped settling a worker from listing absence, a local worker whose shell ended while Orca was closed was never settled: it stayed dispatched on every sweep until `worker-abandon`. Give the daemon the relay's rule instead of a second record of the same fact. `TerminalHost.reapSession` already disposes the exited session's emulator and drops it from `sessions`; it now drops it only when `broadcastExit` reached an attached client. With none attached nobody received the exit, so the emulator-free record stays, and `inspectProcess` hands it over as `verdict: 'exited'` for an exact `expectedIncarnationId` match and deletes it in the same breath. A caller with the wrong incarnation, or none, still gets `SessionNotFoundError`, and a session recreated under the same id replaces the record rather than inheriting it. Every other reader of `sessions` already gates on `isAlive`, so listing, kill, snapshots, sizes and liveness probes are unchanged. The daemon still retires exactly when it did before: `onSessionReaped` fires on exit whether or not the record was dropped. Both daemon-fronting providers were answering `terminal_gone` from their own routing tables for exactly the ids this question is about, and the degraded one dropped `expectedIncarnationId` and `steadyState` entirely. An unclaimed id now routes to the current daemon when — and only when — the caller names a remembered incarnation; a daemon that never held it answers not-found, so routing cannot manufacture an exit, and an ordinary poll still may not borrow a route it never owned. Deletes `retired-pty-incarnations.ts`, its test and the TTL. Refs docs/reference/ssh-execution-boundary.md
This commit is contained in:
@@ -0,0 +1,113 @@
|
||||
// A worker whose shell ended while Orca was closed leaves no route behind: nothing in this
|
||||
// process claims its id any more. The incarnation-scoped question is exactly the one that has to
|
||||
// survive that, so both daemon-fronting providers must carry it (and its `expectedIncarnationId`)
|
||||
// to a daemon instead of answering `terminal_gone` from their own bookkeeping.
|
||||
import { describe, expect, it, vi } from 'vitest'
|
||||
import { DaemonPtyRouter } from './daemon-pty-router'
|
||||
import { DegradedDaemonPtyProvider } from './degraded-daemon-pty-provider'
|
||||
import type { DaemonPtyAdapter } from './daemon-pty-adapter'
|
||||
import type { IPtyProvider } from '../providers/types'
|
||||
import type { PtyProcessInspection } from '../providers/pty-process-inspection'
|
||||
|
||||
const INCARNATION = 'c0ffee00-0000-4000-8000-000000000001'
|
||||
|
||||
const EXITED: PtyProcessInspection = {
|
||||
foregroundProcess: null,
|
||||
hasChildProcesses: false,
|
||||
foregroundProcessEvidence: {
|
||||
authorityGeneration: 'gen-1',
|
||||
observationEpoch: 1,
|
||||
capturedAgeMs: 0,
|
||||
ptyId: 'pty-away',
|
||||
ptyIncarnationId: INCARNATION,
|
||||
verdict: 'exited',
|
||||
reason: 'pty_exit_0'
|
||||
}
|
||||
}
|
||||
|
||||
function createAdapter(sessions: string[] = []): DaemonPtyAdapter {
|
||||
return {
|
||||
hasPty: vi.fn((id: string) => sessions.includes(id)),
|
||||
inspectProcess: vi.fn(async () => EXITED),
|
||||
listProcesses: vi.fn(async () => sessions.map((id) => ({ id, cwd: '', title: 'daemon' }))),
|
||||
onData: vi.fn(() => () => {}),
|
||||
onExit: vi.fn(() => () => {}),
|
||||
onWriteUnavailable: vi.fn(() => () => {}),
|
||||
onBackgroundStreamEvent: vi.fn(() => () => {})
|
||||
} as unknown as DaemonPtyAdapter
|
||||
}
|
||||
|
||||
function createFallbackProvider(): IPtyProvider {
|
||||
return {
|
||||
hasPty: vi.fn(() => false),
|
||||
inspectProcess: vi.fn(async () => EXITED),
|
||||
listProcesses: vi.fn(async () => []),
|
||||
onData: vi.fn(() => () => {}),
|
||||
onExit: vi.fn(() => () => {}),
|
||||
onWriteUnavailable: vi.fn(() => () => {}),
|
||||
onBackgroundStreamEvent: vi.fn(() => () => {})
|
||||
} as unknown as IPtyProvider
|
||||
}
|
||||
|
||||
describe('DaemonPtyRouter incarnation-scoped inspection', () => {
|
||||
it('asks the current daemon about an unclaimed id when the caller names an incarnation', async () => {
|
||||
const current = createAdapter()
|
||||
const router = new DaemonPtyRouter({ current, legacy: [createAdapter()] })
|
||||
|
||||
await expect(
|
||||
router.inspectProcess('pty-away', { expectedIncarnationId: INCARNATION })
|
||||
).resolves.toEqual(EXITED)
|
||||
expect(current.inspectProcess).toHaveBeenCalledWith('pty-away', {
|
||||
expectedIncarnationId: INCARNATION
|
||||
})
|
||||
})
|
||||
|
||||
it('forwards the incarnation to the daemon that still claims the id', async () => {
|
||||
const legacy = createAdapter(['pty-live'])
|
||||
const router = new DaemonPtyRouter({ current: createAdapter(), legacy: [legacy] })
|
||||
|
||||
await router.inspectProcess('pty-live', { expectedIncarnationId: INCARNATION })
|
||||
|
||||
expect(legacy.inspectProcess).toHaveBeenCalledWith('pty-live', {
|
||||
expectedIncarnationId: INCARNATION
|
||||
})
|
||||
})
|
||||
})
|
||||
|
||||
describe('DegradedDaemonPtyProvider incarnation-scoped inspection', () => {
|
||||
it('asks the current daemon about an unclaimed id when the caller names an incarnation', async () => {
|
||||
const current = createAdapter()
|
||||
const provider = new DegradedDaemonPtyProvider({
|
||||
current,
|
||||
legacy: [],
|
||||
fallback: createFallbackProvider()
|
||||
})
|
||||
|
||||
await expect(
|
||||
provider.inspectProcess('pty-away', { expectedIncarnationId: INCARNATION })
|
||||
).resolves.toEqual(EXITED)
|
||||
expect(current.inspectProcess).toHaveBeenCalledWith('pty-away', {
|
||||
expectedIncarnationId: INCARNATION
|
||||
})
|
||||
})
|
||||
|
||||
it('forwards inspection options to the provider that owns the session', async () => {
|
||||
const current = createAdapter(['pty-live'])
|
||||
const provider = new DegradedDaemonPtyProvider({
|
||||
current,
|
||||
legacy: [],
|
||||
fallback: createFallbackProvider()
|
||||
})
|
||||
await provider.discoverDaemonSessions()
|
||||
|
||||
await provider.inspectProcess('pty-live', {
|
||||
expectedIncarnationId: INCARNATION,
|
||||
steadyState: true
|
||||
})
|
||||
|
||||
expect(current.inspectProcess).toHaveBeenCalledWith('pty-live', {
|
||||
expectedIncarnationId: INCARNATION,
|
||||
steadyState: true
|
||||
})
|
||||
})
|
||||
})
|
||||
@@ -182,7 +182,7 @@ export class DaemonPtyRouter implements IPtyProvider {
|
||||
id: string,
|
||||
options?: { expectedIncarnationId?: string; steadyState?: boolean }
|
||||
): Promise<PtyProcessInspection> {
|
||||
return this.adapterForInspection(id).inspectProcess(id, options)
|
||||
return this.adapterForInspection(id, options?.expectedIncarnationId).inspectProcess(id, options)
|
||||
}
|
||||
|
||||
async confirmForegroundProcess(id: string): Promise<string | null> {
|
||||
@@ -325,15 +325,25 @@ export class DaemonPtyRouter implements IPtyProvider {
|
||||
return this.sessionAdapters.get(sessionId) ?? this.current
|
||||
}
|
||||
|
||||
private adapterForInspection(sessionId: string): DaemonPtyAdapter {
|
||||
private adapterForInspection(
|
||||
sessionId: string,
|
||||
expectedIncarnationId?: string
|
||||
): DaemonPtyAdapter {
|
||||
const adapter =
|
||||
this.sessionAdapters.get(sessionId) ??
|
||||
this.allAdapters().find((candidate) => candidate.hasPty(sessionId))
|
||||
if (!adapter) {
|
||||
if (adapter) {
|
||||
this.sessionAdapters.set(sessionId, adapter)
|
||||
return adapter
|
||||
}
|
||||
// An unclaimed id is exactly what a caller naming a remembered incarnation asks about: the
|
||||
// session that died while this client was away, so no route survived. A daemon that never held
|
||||
// it answers not-found, so routing the question cannot manufacture an exit. A caller that names
|
||||
// no incarnation is an ordinary poll and must not borrow a route it never owned.
|
||||
if (expectedIncarnationId === undefined) {
|
||||
throw new Error('terminal_gone')
|
||||
}
|
||||
this.sessionAdapters.set(sessionId, adapter)
|
||||
return adapter
|
||||
return this.current
|
||||
}
|
||||
|
||||
private allAdapters(): DaemonPtyAdapter[] {
|
||||
|
||||
@@ -1,7 +1,7 @@
|
||||
import type { DaemonPtyAdapter } from './daemon-pty-adapter'
|
||||
import { combineUnsubscribes } from './combine-unsubscribes'
|
||||
import { shutdownDegradedFallbackSessions } from './degraded-daemon-fallback-shutdown'
|
||||
import { inspectPtyProviderProcess } from '../providers/pty-process-inspection'
|
||||
import type { PtyProcessInspectionOptions } from '../providers/pty-process-inspection'
|
||||
import type {
|
||||
IPtyProvider,
|
||||
PtyBackgroundStreamEvent,
|
||||
@@ -15,6 +15,7 @@ import {
|
||||
adoptOwningProvider,
|
||||
attachDaemonOwnedSession,
|
||||
findDaemonAdapter,
|
||||
inspectRoutedDaemonProcess,
|
||||
listProviderSessionIds
|
||||
} from './degraded-daemon-session-routing'
|
||||
import { DegradedDaemonFreshSpawnRouter } from './degraded-daemon-fresh-spawn-routing'
|
||||
@@ -184,10 +185,9 @@ export class DegradedDaemonPtyProvider implements IPtyProvider {
|
||||
async getForegroundProcess(id: string): Promise<string | null> {
|
||||
return this.providerFor(id).getForegroundProcess(id)
|
||||
}
|
||||
inspectProcess(id: string) {
|
||||
return this.hasPty(id)
|
||||
? inspectPtyProviderProcess(this.providerFor(id), id)
|
||||
: Promise.reject(new Error('terminal_gone'))
|
||||
inspectProcess(id: string, options?: PtyProcessInspectionOptions) {
|
||||
const routed = this.hasPty(id) ? this.providerFor(id) : null
|
||||
return inspectRoutedDaemonProcess(routed, this.current, id, options)
|
||||
}
|
||||
async confirmForegroundProcess(id: string): Promise<string | null> {
|
||||
return this.providerFor(id).confirmForegroundProcess?.(id) ?? null
|
||||
|
||||
@@ -1,6 +1,31 @@
|
||||
import type { IPtyProvider } from '../providers/types'
|
||||
import type { DaemonPtyAdapter } from './daemon-pty-adapter'
|
||||
import { SessionNotFoundError } from './daemon-errors'
|
||||
import {
|
||||
inspectPtyProviderProcess,
|
||||
type PtyProcessInspection,
|
||||
type PtyProcessInspectionOptions
|
||||
} from '../providers/pty-process-inspection'
|
||||
|
||||
/**
|
||||
* An id no route claims may only be asked about by a caller naming a remembered incarnation: that
|
||||
* is exactly the session that died while this client was away, so no route survived it. A daemon
|
||||
* that never held it answers not-found, so routing the question cannot manufacture an exit. An
|
||||
* ordinary poll must not borrow a route it never owned. Mirrors DaemonPtyRouter.
|
||||
*/
|
||||
export function inspectRoutedDaemonProcess(
|
||||
owner: IPtyProvider | null,
|
||||
currentDaemon: DaemonPtyAdapter,
|
||||
sessionId: string,
|
||||
options?: PtyProcessInspectionOptions
|
||||
): Promise<PtyProcessInspection> {
|
||||
if (owner) {
|
||||
return inspectPtyProviderProcess(owner, sessionId, options)
|
||||
}
|
||||
return options?.expectedIncarnationId === undefined
|
||||
? Promise.reject(new Error('terminal_gone'))
|
||||
: currentDaemon.inspectProcess(sessionId, options)
|
||||
}
|
||||
|
||||
export function listProviderSessionIds(
|
||||
sessionProviders: ReadonlyMap<string, IPtyProvider>,
|
||||
|
||||
@@ -1,30 +0,0 @@
|
||||
import { describe, expect, it } from 'vitest'
|
||||
import { pruneRetiredPtyIncarnations } from './retired-pty-incarnations'
|
||||
|
||||
describe('retired PTY incarnation retention', () => {
|
||||
it('removes expired records before they can accumulate', () => {
|
||||
const records = new Map([
|
||||
['expired', { incarnationId: 'a', code: 0, expiresAt: 10 }],
|
||||
['live', { incarnationId: 'b', code: 0, expiresAt: 30 }]
|
||||
])
|
||||
|
||||
pruneRetiredPtyIncarnations(records, 20)
|
||||
|
||||
expect([...records.keys()]).toEqual(['live'])
|
||||
})
|
||||
|
||||
it('caps records when many distinct PTYs retire together', () => {
|
||||
const records = new Map(
|
||||
Array.from({ length: 1001 }, (_, index) => [
|
||||
`pty-${index}`,
|
||||
{ incarnationId: `inc-${index}`, code: 0, expiresAt: 100 }
|
||||
])
|
||||
)
|
||||
|
||||
pruneRetiredPtyIncarnations(records, 0)
|
||||
|
||||
expect(records.size).toBe(1000)
|
||||
expect(records.has('pty-0')).toBe(false)
|
||||
expect(records.has('pty-1000')).toBe(true)
|
||||
})
|
||||
})
|
||||
@@ -1,26 +0,0 @@
|
||||
export type RetiredPtyIncarnation = {
|
||||
incarnationId: string
|
||||
code: number
|
||||
expiresAt: number
|
||||
}
|
||||
|
||||
const MAX_RETIRED_PTY_INCARNATIONS = 1000
|
||||
|
||||
/** Drop expired exit evidence and cap retained records during long-lived hosts. */
|
||||
export function pruneRetiredPtyIncarnations(
|
||||
records: Map<string, RetiredPtyIncarnation>,
|
||||
now = Date.now()
|
||||
): void {
|
||||
for (const [id, record] of records) {
|
||||
if (record.expiresAt <= now) {
|
||||
records.delete(id)
|
||||
}
|
||||
}
|
||||
while (records.size > MAX_RETIRED_PTY_INCARNATIONS) {
|
||||
const oldest = records.keys().next().value
|
||||
if (oldest === undefined) {
|
||||
break
|
||||
}
|
||||
records.delete(oldest)
|
||||
}
|
||||
}
|
||||
@@ -1,8 +1,11 @@
|
||||
import { describe, expect, it, vi } from 'vitest'
|
||||
import type { SubprocessHandle } from './session-subprocess-handle'
|
||||
import { TerminalHost } from './terminal-host'
|
||||
import { SessionNotFoundError } from './types'
|
||||
|
||||
function createSubprocess(): SubprocessHandle {
|
||||
type MockSubprocess = SubprocessHandle & { exit(code: number): void }
|
||||
|
||||
function createSubprocess(): MockSubprocess {
|
||||
let onExit: ((code: number) => void) | null = null
|
||||
return {
|
||||
pid: 99_999,
|
||||
@@ -17,8 +20,20 @@ function createSubprocess(): SubprocessHandle {
|
||||
onExit: (callback) => {
|
||||
onExit = callback
|
||||
},
|
||||
dispose: vi.fn()
|
||||
}
|
||||
dispose: vi.fn(),
|
||||
exit: (code: number) => onExit?.(code)
|
||||
} as MockSubprocess
|
||||
}
|
||||
|
||||
function createHost(): { host: TerminalHost; lastSubprocess: () => MockSubprocess } {
|
||||
let last: MockSubprocess | undefined
|
||||
const host = new TerminalHost({
|
||||
spawnSubprocess: () => {
|
||||
last = createSubprocess()
|
||||
return last
|
||||
}
|
||||
})
|
||||
return { host, lastSubprocess: () => last as MockSubprocess }
|
||||
}
|
||||
|
||||
describe('TerminalHost process inspection', () => {
|
||||
@@ -47,3 +62,127 @@ describe('TerminalHost process inspection', () => {
|
||||
}
|
||||
})
|
||||
})
|
||||
|
||||
describe('TerminalHost undelivered exits', () => {
|
||||
/** Orca quits (the client's attachments drop), the shell ends, Orca reopens and asks. */
|
||||
async function exitWhileClientIsAway(
|
||||
host: TerminalHost,
|
||||
subprocess: () => MockSubprocess,
|
||||
sessionId: string,
|
||||
code = 0
|
||||
): Promise<string> {
|
||||
const created = await host.createOrAttach({
|
||||
sessionId,
|
||||
cols: 80,
|
||||
rows: 24,
|
||||
streamClient: { onData: vi.fn(), onExit: vi.fn() }
|
||||
})
|
||||
host.detach(sessionId, created.attachToken as symbol)
|
||||
subprocess().exit(code)
|
||||
return created.incarnationId
|
||||
}
|
||||
|
||||
it('hands a close-and-reopen caller the exit its client never received, once', async () => {
|
||||
const { host, lastSubprocess } = createHost()
|
||||
try {
|
||||
const incarnationId = await exitWhileClientIsAway(host, lastSubprocess, 'session-away', 3)
|
||||
|
||||
// The exit is gone from every liveness surface, but the evidence is still owed to its owner.
|
||||
expect(host.listSessions()).toHaveLength(0)
|
||||
await expect(
|
||||
host.inspectProcess('session-away', { expectedIncarnationId: incarnationId })
|
||||
).resolves.toMatchObject({
|
||||
foregroundProcess: null,
|
||||
hasChildProcesses: false,
|
||||
foregroundProcessEvidence: {
|
||||
ptyId: 'session-away',
|
||||
ptyIncarnationId: incarnationId,
|
||||
verdict: 'exited',
|
||||
reason: 'pty_exit_3'
|
||||
}
|
||||
})
|
||||
|
||||
// Handing it over was the delivery: that client now has the exit.
|
||||
expect(() =>
|
||||
host.inspectProcess('session-away', { expectedIncarnationId: incarnationId })
|
||||
).toThrow(SessionNotFoundError)
|
||||
} finally {
|
||||
await host.dispose()
|
||||
}
|
||||
})
|
||||
|
||||
it('answers nothing to a caller with the wrong incarnation or none at all', async () => {
|
||||
const { host, lastSubprocess } = createHost()
|
||||
try {
|
||||
const incarnationId = await exitWhileClientIsAway(host, lastSubprocess, 'session-away')
|
||||
|
||||
expect(() =>
|
||||
host.inspectProcess('session-away', { expectedIncarnationId: 'someone-else' })
|
||||
).toThrow(SessionNotFoundError)
|
||||
expect(() => host.inspectProcess('session-away')).toThrow(SessionNotFoundError)
|
||||
|
||||
// Neither refusal consumed the evidence the real owner is still owed.
|
||||
await expect(
|
||||
host.inspectProcess('session-away', { expectedIncarnationId: incarnationId })
|
||||
).resolves.toMatchObject({ foregroundProcessEvidence: { verdict: 'exited' } })
|
||||
} finally {
|
||||
await host.dispose()
|
||||
}
|
||||
})
|
||||
|
||||
it('holds nothing when the exit reached an attached client', async () => {
|
||||
const { host, lastSubprocess } = createHost()
|
||||
try {
|
||||
const onExit = vi.fn()
|
||||
const created = await host.createOrAttach({
|
||||
sessionId: 'session-attached',
|
||||
cols: 80,
|
||||
rows: 24,
|
||||
streamClient: { onData: vi.fn(), onExit }
|
||||
})
|
||||
lastSubprocess().exit(0)
|
||||
|
||||
expect(onExit).toHaveBeenCalledWith(0, created.incarnationId, expect.anything())
|
||||
expect(() =>
|
||||
host.inspectProcess('session-attached', { expectedIncarnationId: created.incarnationId })
|
||||
).toThrow(SessionNotFoundError)
|
||||
} finally {
|
||||
await host.dispose()
|
||||
}
|
||||
})
|
||||
|
||||
it('does not let a session recreated under the same id inherit the old exit', async () => {
|
||||
const { host, lastSubprocess } = createHost()
|
||||
try {
|
||||
const incarnationId = await exitWhileClientIsAway(host, lastSubprocess, 'session-reused')
|
||||
const replacement = await host.createOrAttach({
|
||||
sessionId: 'session-reused',
|
||||
cols: 80,
|
||||
rows: 24,
|
||||
streamClient: { onData: vi.fn(), onExit: vi.fn() }
|
||||
})
|
||||
|
||||
expect(replacement.incarnationId).not.toBe(incarnationId)
|
||||
await expect(
|
||||
host.inspectProcess('session-reused', { expectedIncarnationId: incarnationId })
|
||||
).resolves.toMatchObject({
|
||||
foregroundProcessEvidence: { verdict: 'unverifiable', reason: 'incarnation_mismatch' }
|
||||
})
|
||||
} finally {
|
||||
await host.dispose()
|
||||
}
|
||||
})
|
||||
|
||||
it('keeps kill tombstones independent of a held exit', async () => {
|
||||
const { host, lastSubprocess } = createHost()
|
||||
try {
|
||||
await exitWhileClientIsAway(host, lastSubprocess, 'session-away')
|
||||
|
||||
// Nothing killed it, and a held record is not a live session to kill.
|
||||
expect(host.isKilled('session-away')).toBe(false)
|
||||
expect(() => host.kill('session-away')).toThrow(SessionNotFoundError)
|
||||
} finally {
|
||||
await host.dispose()
|
||||
}
|
||||
})
|
||||
})
|
||||
|
||||
@@ -19,7 +19,8 @@ export type TerminalHostProcessInspection = {
|
||||
foregroundProcessEvidence?: RemoteForegroundEvidence
|
||||
}
|
||||
|
||||
type RetiredIncarnation = { incarnationId: string; code: number; expiresAt: number }
|
||||
/** An observed exit the session's owner never received; the host matched the caller's incarnation. */
|
||||
type UndeliveredExit = { incarnationId: string; code: number }
|
||||
|
||||
/**
|
||||
* Tick tiers for a POSIX pane. `cheap` forks `ps` without `tty=`/`command=` (11-38x cheaper)
|
||||
@@ -34,18 +35,14 @@ export async function inspectTerminalHostProcess(args: {
|
||||
expectedIncarnationId?: string
|
||||
/** The caller is a self-correcting poll that only reads the process name, never evidence. */
|
||||
steadyState?: boolean
|
||||
retiredIncarnation?: RetiredIncarnation
|
||||
undeliveredExit?: UndeliveredExit
|
||||
authorityGeneration: string
|
||||
nextObservationEpoch: () => number
|
||||
onTier?: (tier: TerminalHostInspectionTier) => void
|
||||
}): Promise<TerminalHostProcessInspection> {
|
||||
const { sessionId, session, expectedIncarnationId, retiredIncarnation } = args
|
||||
const { sessionId, session, expectedIncarnationId, undeliveredExit } = args
|
||||
if (!session || !session.isAlive) {
|
||||
if (
|
||||
retiredIncarnation &&
|
||||
retiredIncarnation.expiresAt > Date.now() &&
|
||||
expectedIncarnationId === retiredIncarnation.incarnationId
|
||||
) {
|
||||
if (undeliveredExit) {
|
||||
return {
|
||||
foregroundProcess: null,
|
||||
hasChildProcesses: false,
|
||||
@@ -54,9 +51,9 @@ export async function inspectTerminalHostProcess(args: {
|
||||
observationEpoch: args.nextObservationEpoch(),
|
||||
capturedAgeMs: 0,
|
||||
ptyId: sessionId,
|
||||
ptyIncarnationId: retiredIncarnation.incarnationId,
|
||||
ptyIncarnationId: undeliveredExit.incarnationId,
|
||||
verdict: 'exited',
|
||||
reason: `pty_exit_${retiredIncarnation.code}`
|
||||
reason: `pty_exit_${undeliveredExit.code}`
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
@@ -22,7 +22,6 @@ import { createOrAttachTerminalSession } from './terminal-host-session-create'
|
||||
import { TerminalAttachCanceledError } from './daemon-errors'
|
||||
import { rejectOnAbort } from './terminal-attach-cancellation'
|
||||
import { randomUUID } from 'node:crypto'
|
||||
import { pruneRetiredPtyIncarnations } from './retired-pty-incarnations'
|
||||
import {
|
||||
inspectTerminalHostProcess,
|
||||
type TerminalHostProcessInspection
|
||||
@@ -42,7 +41,6 @@ export type { CreateOrAttachOptions, CreateOrAttachResult } from './terminal-hos
|
||||
export type { TerminalHostOptions } from './terminal-host-options'
|
||||
|
||||
const DEFAULT_MAX_TOMBSTONES = 1000
|
||||
const REMOTE_FOREGROUND_TOMBSTONE_RETENTION_MS = 2_000
|
||||
|
||||
export class TerminalHost {
|
||||
private sessions = new Map<string, Session>()
|
||||
@@ -61,10 +59,6 @@ export class TerminalHost {
|
||||
private readonly agentSessionGenerations = new TerminalHostAgentSessionGenerations()
|
||||
private readonly authorityGeneration = randomUUID()
|
||||
private observationEpoch = 0
|
||||
private readonly retiredIncarnations = new Map<
|
||||
string,
|
||||
{ incarnationId: string; code: number; expiresAt: number }
|
||||
>()
|
||||
|
||||
constructor(opts: TerminalHostOptions) {
|
||||
this.spawnSubprocess = opts.spawnSubprocess
|
||||
@@ -124,15 +118,6 @@ export class TerminalHost {
|
||||
? { reportReadinessEvent: this.reportReadinessEvent }
|
||||
: {}),
|
||||
onSessionExit: (sessionId, generation) => {
|
||||
const session = this.sessions.get(sessionId)
|
||||
if (session) {
|
||||
pruneRetiredPtyIncarnations(this.retiredIncarnations)
|
||||
this.retiredIncarnations.set(sessionId, {
|
||||
incarnationId: session.incarnationId,
|
||||
code: session.exitCode ?? 0,
|
||||
expiresAt: Date.now() + REMOTE_FOREGROUND_TOMBSTONE_RETENTION_MS
|
||||
})
|
||||
}
|
||||
this.agentSessionOwners.release(sessionId, generation)
|
||||
this.agentSessionGenerations.forget(sessionId, generation)
|
||||
this.reapSession(sessionId)
|
||||
@@ -199,8 +184,14 @@ export class TerminalHost {
|
||||
if (!session || session.isAlive) {
|
||||
return
|
||||
}
|
||||
// `broadcastExit` just fanned this exit out to the clients attached at that instant. With none
|
||||
// attached nobody received it, so the (now emulator-free) record stays until a caller naming
|
||||
// this exact incarnation is handed it -- see takeUndeliveredExit.
|
||||
const delivered = session.hasAttachedClients
|
||||
session.dispose()
|
||||
this.sessions.delete(sessionId)
|
||||
if (delivered) {
|
||||
this.sessions.delete(sessionId)
|
||||
}
|
||||
this.onSessionReaped?.(sessionId)
|
||||
}
|
||||
|
||||
@@ -235,15 +226,11 @@ export class TerminalHost {
|
||||
sessionId: string,
|
||||
options?: { expectedIncarnationId?: string; steadyState?: boolean }
|
||||
): Promise<TerminalHostProcessInspection> {
|
||||
pruneRetiredPtyIncarnations(this.retiredIncarnations)
|
||||
const session = this.sessions.get(sessionId)
|
||||
if (
|
||||
(!session || !session.isAlive) &&
|
||||
!(
|
||||
(this.retiredIncarnations.get(sessionId)?.expiresAt ?? 0) > Date.now() &&
|
||||
options?.expectedIncarnationId === this.retiredIncarnations.get(sessionId)?.incarnationId
|
||||
)
|
||||
) {
|
||||
const undeliveredExit = session?.isAlive
|
||||
? undefined
|
||||
: this.takeUndeliveredExit(sessionId, options?.expectedIncarnationId)
|
||||
if (!session?.isAlive && !undeliveredExit) {
|
||||
// Preserve the historical synchronous missing-session failure.
|
||||
throw new SessionNotFoundError(sessionId)
|
||||
}
|
||||
@@ -254,12 +241,36 @@ export class TerminalHost {
|
||||
? { expectedIncarnationId: options.expectedIncarnationId }
|
||||
: {}),
|
||||
...(options?.steadyState === true ? { steadyState: true } : {}),
|
||||
retiredIncarnation: this.retiredIncarnations.get(sessionId),
|
||||
...(undeliveredExit ? { undeliveredExit } : {}),
|
||||
authorityGeneration: this.authorityGeneration,
|
||||
nextObservationEpoch: () => ++this.observationEpoch
|
||||
})
|
||||
}
|
||||
|
||||
/**
|
||||
* The exit this session's owner never received, for the incarnation the caller named. Handing it
|
||||
* over IS the delivery, so the record leaves with it and a later ask reads as not-found -- that
|
||||
* caller already has the exit. A caller that names no incarnation, or a different one, proves
|
||||
* nothing about this process and gets the ordinary missing-session failure
|
||||
* (docs/reference/ssh-execution-boundary.md).
|
||||
*/
|
||||
private takeUndeliveredExit(
|
||||
sessionId: string,
|
||||
expectedIncarnationId: string | undefined
|
||||
): { incarnationId: string; code: number } | undefined {
|
||||
const session = this.sessions.get(sessionId)
|
||||
if (
|
||||
!session ||
|
||||
session.isAlive ||
|
||||
session.exitCode === null ||
|
||||
session.incarnationId !== expectedIncarnationId
|
||||
) {
|
||||
return undefined
|
||||
}
|
||||
this.sessions.delete(sessionId)
|
||||
return { incarnationId: session.incarnationId, code: session.exitCode }
|
||||
}
|
||||
|
||||
async confirmForegroundProcess(sessionId: string): Promise<string | null> {
|
||||
return confirmTerminalHostForegroundProcess(this.sessions.get(sessionId))
|
||||
}
|
||||
|
||||
+129
@@ -16,6 +16,48 @@ import {
|
||||
store
|
||||
} from '../orca-runtime-test-fixtures.spec'
|
||||
import { publishLegacyWorkerReveal } from '../orca-runtime-test-scenario-builders.spec'
|
||||
import { TerminalHost } from '../../daemon/terminal-host'
|
||||
import { DaemonPtyRouter } from '../../daemon/daemon-pty-router'
|
||||
import type { DaemonPtyAdapter } from '../../daemon/daemon-pty-adapter'
|
||||
import type { SubprocessHandle } from '../../daemon/session-subprocess-handle'
|
||||
import { getLocalPtyProvider, setLocalPtyProvider } from '../../ipc/pty/provider/registry'
|
||||
import { inspectExitedIncarnationFromRuntimeController } from '../../ipc/pty/runtime/operations'
|
||||
|
||||
/** A daemon that outlived the app close, fronted by the router main actually asks. */
|
||||
function daemonRouterOver(host: TerminalHost): DaemonPtyRouter {
|
||||
const adapter = {
|
||||
// Nothing in this process routes the id any more: the app restarted after the shell ended.
|
||||
hasPty: () => false,
|
||||
inspectProcess: (id: string, options?: { expectedIncarnationId?: string }) =>
|
||||
host.inspectProcess(id, options),
|
||||
listProcesses: async () => [],
|
||||
onData: () => () => {},
|
||||
onExit: () => () => {},
|
||||
onWriteUnavailable: () => () => {},
|
||||
onBackgroundStreamEvent: () => () => {}
|
||||
} as unknown as DaemonPtyAdapter
|
||||
return new DaemonPtyRouter({ current: adapter, legacy: [] })
|
||||
}
|
||||
|
||||
function daemonSubprocess(): SubprocessHandle & { exit(code: number): void } {
|
||||
let onExit: ((code: number) => void) | null = null
|
||||
return {
|
||||
pid: 4242,
|
||||
getForegroundProcess: () => null,
|
||||
write: vi.fn(),
|
||||
resize: vi.fn(),
|
||||
kill: vi.fn(),
|
||||
terminateOwnedTree: () => 'unavailable',
|
||||
forceKill: vi.fn(),
|
||||
signal: vi.fn(),
|
||||
onData: vi.fn(),
|
||||
onExit: (callback: (code: number) => void) => {
|
||||
onExit = callback
|
||||
},
|
||||
dispose: vi.fn(),
|
||||
exit: (code: number) => onExit?.(code)
|
||||
} as unknown as SubprocessHandle & { exit(code: number): void }
|
||||
}
|
||||
|
||||
describe('OrcaRuntimeService', () => {
|
||||
it('requeues an active Task before clearing recovery for an authoritatively missing worker', async () => {
|
||||
@@ -199,6 +241,93 @@ describe('OrcaRuntimeService', () => {
|
||||
}
|
||||
})
|
||||
|
||||
it('settles a local worker whose daemon session exited while the app was closed', async () => {
|
||||
const workerPaneKey = `legacy-daemon-exit:${HEADLESS_LEAF_ID}`
|
||||
const ptyId = 'pty-daemon-worker'
|
||||
const subprocess = daemonSubprocess()
|
||||
const host = new TerminalHost({ spawnSubprocess: () => subprocess })
|
||||
// Orca quits (the client's attachment drops), then the worker's shell ends.
|
||||
const created = await host.createOrAttach({
|
||||
sessionId: ptyId,
|
||||
cols: 80,
|
||||
rows: 24,
|
||||
streamClient: { onData: vi.fn(), onExit: vi.fn() }
|
||||
})
|
||||
host.detach(ptyId, created.attachToken as symbol)
|
||||
subprocess.exit(0)
|
||||
|
||||
const { runtimeStore, getSession } = makeRuntimeStoreWithWorkspaceSession({
|
||||
...getDefaultWorkspaceSession(),
|
||||
tabsByWorktree: { [TEST_WORKTREE_ID]: [] },
|
||||
sleepingAgentSessionsByPaneKey: {
|
||||
[workerPaneKey]: {
|
||||
paneKey: workerPaneKey,
|
||||
tabId: 'legacy-daemon-exit',
|
||||
worktreeId: TEST_WORKTREE_ID,
|
||||
agent: 'codex',
|
||||
providerSession: { key: 'session_id', id: 'legacy-daemon-exit-session' },
|
||||
prompt: 'continue',
|
||||
state: 'working',
|
||||
capturedAt: 1,
|
||||
updatedAt: 1,
|
||||
origin: 'live'
|
||||
}
|
||||
}
|
||||
})
|
||||
const runtime = new OrcaRuntimeService(
|
||||
{ ...runtimeStore, flushOrThrow: vi.fn() } as never,
|
||||
undefined,
|
||||
{ canRecoverPersistentLocalPtys: () => true }
|
||||
)
|
||||
const db = new OrchestrationDb(':memory:')
|
||||
const previousProvider = getLocalPtyProvider()
|
||||
setLocalPtyProvider(daemonRouterOver(host))
|
||||
try {
|
||||
const task = db.createTask({ runId: 'run_legacy_local', spec: 'daemon worker' })
|
||||
const started = db.createStartingWorkerDispatch({
|
||||
creator: { kind: 'system' },
|
||||
maxDepth: Number.MAX_SAFE_INTEGER,
|
||||
taskId: task.id,
|
||||
startOptions: { topology: 'current', agent: 'codex' }
|
||||
})
|
||||
db.prepareStartingWorkerAuthority({
|
||||
dispatchId: started.dispatch.id,
|
||||
handle: 'term_daemon_exit',
|
||||
paneKey: workerPaneKey,
|
||||
processIncarnation: `${ptyId}:${created.incarnationId}`,
|
||||
worktreeId: TEST_WORKTREE_ID,
|
||||
setupState: 'not_applicable',
|
||||
effects: []
|
||||
})
|
||||
db.markWorkerDispatchReady(started.dispatch.id)
|
||||
runtime.setOrchestrationDb(db)
|
||||
// No stub: this is the shipped proof chain, from the recovery port through
|
||||
// providerObservedIncarnationExit and the daemon router to the host that watched it die.
|
||||
runtime.setPtyController({
|
||||
write: vi.fn(() => true),
|
||||
kill: vi.fn(() => true),
|
||||
getForegroundProcess: async () => null,
|
||||
hasPty: () => false,
|
||||
inspectExitedIncarnation: (candidatePtyId, incarnationId) =>
|
||||
inspectExitedIncarnationFromRuntimeController(candidatePtyId, incarnationId),
|
||||
listProcesses: async () => []
|
||||
})
|
||||
runtime.setNotifier({ resolveLegacyWorkerTerminalRecovery: vi.fn() } as never)
|
||||
|
||||
await expect(runtime.reconcileLegacyWorkerTerminals()).resolves.toMatchObject({
|
||||
adoptedDispatchIds: [],
|
||||
exitedDispatchIds: [started.dispatch.id],
|
||||
deferredDispatchIds: []
|
||||
})
|
||||
expect(db.getDispatchContextById(started.dispatch.id)?.status).not.toBe('dispatched')
|
||||
expect(getSession().sleepingAgentSessionsByPaneKey?.[workerPaneKey]).toBeUndefined()
|
||||
} finally {
|
||||
setLocalPtyProvider(previousProvider)
|
||||
await host.dispose()
|
||||
db.close()
|
||||
}
|
||||
})
|
||||
|
||||
it('waits for durability before retry settles a resolution already present in memory', async () => {
|
||||
const workerPaneKey = `legacy-missing-retry:${HEADLESS_LEAF_ID}`
|
||||
const incarnationId = '34343434-3434-4434-8434-343434343434'
|
||||
|
||||
Reference in New Issue
Block a user