mirror of
https://github.com/stablyai/orca.git
synced 2026-10-02 16:02:15 +00:00
* feat(daemon): idle retirement, session census and recovery-only provider (#16741 T2 P4b)
- DaemonPtyRouter routes idle retirement through DaemonRouterRetirement: it
fences new spawns, counts spawns already in flight, takes a census of every
daemon generation and retires them only when all are idle. A lost reply
keeps the fence; an unanswered census is unverifiable, never empty.
- Each adapter answers requestIdleRetirement through shutdownIfIdle and fences
its own spawns while retiring.
- A recovery-only adapter/provider (createDaemonRecoveryProvider) reattaches
and controls existing sessions but admits no new process, never respawns
the daemon, and never prunes sessions of unknown worktrees.
- listLiveDaemonSessions and requestIdleDaemonRetirement report null /
'unverifiable' when any generation cannot answer.
- ptySpawnHealth replies add optional coverage and Node runtime fields;
checkDaemonHealthWithCoverage reads them and treats an absent coverage from
an older Windows daemon as handshake-only.
- reconcileOnStartup moves to daemon-router-session-reconciliation unchanged.
No caller uses retirement, the census or the recovery provider yet (T6), so
they are inert. Ported by hunk from #16741 (a68b6f3531), without the PTY
ownership-transfer input fence (T7) and with a Node-only runtime report.
* test(daemon): narrow parsed frames and endpoint errors instead of asserting types
---------
Co-authored-by: m4air <m4air@m4airs-Air.localdomain>
This commit is contained in:
@@ -0,0 +1,13 @@
|
||||
/** What a `ptySpawnHealth` reply proves about this daemon; every field is optional on the wire. */
|
||||
export function readDaemonHealthIdentity(): {
|
||||
coverage: 'pty-spawn' | 'handshake'
|
||||
runtimeKind: 'node'
|
||||
runtimeVersion: string
|
||||
} {
|
||||
return {
|
||||
// Why handshake on Windows: preflightPtySpawnHealth skips the spawn probe there.
|
||||
coverage: process.platform === 'win32' ? 'handshake' : 'pty-spawn',
|
||||
runtimeKind: 'node',
|
||||
runtimeVersion: process.version
|
||||
}
|
||||
}
|
||||
@@ -11,6 +11,7 @@ import { getDaemonPidPath, serializeDaemonPidFile } from './daemon-spawner'
|
||||
import type { SocketProbeOutcome } from './daemon-endpoint-probe'
|
||||
import {
|
||||
checkDaemonHealth,
|
||||
checkDaemonHealthWithCoverage,
|
||||
E2E_FORCE_DAEMON_HEALTH_UNREACHABLE_ENV,
|
||||
healthCheckDaemon
|
||||
} from './daemon-health'
|
||||
@@ -105,12 +106,63 @@ describe('daemon health', () => {
|
||||
try {
|
||||
await expect(checkDaemonHealth(socketPath, tokenPath)).resolves.toBe('healthy')
|
||||
await expect(healthCheckDaemon(socketPath, tokenPath)).resolves.toBe(true)
|
||||
expect(ptySpawnHealthCheck).toHaveBeenCalledTimes(2)
|
||||
await expect(checkDaemonHealthWithCoverage(socketPath, tokenPath)).resolves.toMatchObject({
|
||||
verdict: 'healthy',
|
||||
coverage: process.platform === 'win32' ? 'handshake' : 'pty-spawn',
|
||||
runtimeKind: 'node',
|
||||
runtimeVersion: process.version
|
||||
})
|
||||
expect(ptySpawnHealthCheck).toHaveBeenCalledTimes(3)
|
||||
} finally {
|
||||
await server.shutdown()
|
||||
}
|
||||
})
|
||||
|
||||
it('treats missing coverage from a legacy Windows daemon as handshake-only', async () => {
|
||||
writeFileSync(tokenPath, 'legacy-token')
|
||||
const server = createServer((socket) => {
|
||||
let pending = ''
|
||||
socket.on('data', (chunk) => {
|
||||
pending += chunk.toString()
|
||||
for (;;) {
|
||||
const newline = pending.indexOf('\n')
|
||||
if (newline === -1) {
|
||||
return
|
||||
}
|
||||
const message: unknown = JSON.parse(pending.slice(0, newline))
|
||||
const type =
|
||||
typeof message === 'object' && message !== null && 'type' in message
|
||||
? message.type
|
||||
: undefined
|
||||
pending = pending.slice(newline + 1)
|
||||
if (type === 'hello') {
|
||||
socket.write(`${JSON.stringify({ type: 'hello', ok: true })}\n`)
|
||||
} else if (type === 'ptySpawnHealth') {
|
||||
socket.write(
|
||||
`${JSON.stringify({ id: 'health-1', ok: true, payload: { healthy: true } })}\n`
|
||||
)
|
||||
}
|
||||
}
|
||||
})
|
||||
})
|
||||
await new Promise<void>((resolve, reject) => {
|
||||
server.once('error', reject)
|
||||
server.listen(socketPath, resolve)
|
||||
})
|
||||
|
||||
const platform = Object.getOwnPropertyDescriptor(process, 'platform')!
|
||||
Object.defineProperty(process, 'platform', { configurable: true, value: 'win32' })
|
||||
try {
|
||||
await expect(checkDaemonHealthWithCoverage(socketPath, tokenPath)).resolves.toEqual({
|
||||
verdict: 'healthy',
|
||||
coverage: 'handshake'
|
||||
})
|
||||
} finally {
|
||||
Object.defineProperty(process, 'platform', platform)
|
||||
await closeServer(server)
|
||||
}
|
||||
})
|
||||
|
||||
it('fails when a protocol-healthy daemon cannot spawn PTYs', async () => {
|
||||
const server = new DaemonServer({
|
||||
socketPath,
|
||||
|
||||
@@ -21,15 +21,57 @@ export const E2E_FORCE_DAEMON_HEALTH_UNREACHABLE_ENV = 'ORCA_E2E_FORCE_DAEMON_HE
|
||||
// also covers a live-but-wedged daemon that simply missed the RPC budget.
|
||||
export type DaemonHealth = 'healthy' | 'unreachable' | 'rejected' | 'pty-spawn-unhealthy'
|
||||
|
||||
export function checkDaemonHealth(socketPath: string, tokenPath: string): Promise<DaemonHealth> {
|
||||
export type DaemonHealthCheck = {
|
||||
verdict: DaemonHealth
|
||||
coverage: 'pty-spawn' | 'handshake'
|
||||
/** Optional runtime proof from newer daemons; absent on mixed-version peers. */
|
||||
runtimeKind?: 'node'
|
||||
runtimeVersion?: string
|
||||
}
|
||||
|
||||
function readRuntimeIdentity(
|
||||
payload: unknown
|
||||
): Pick<DaemonHealthCheck, 'runtimeKind' | 'runtimeVersion'> {
|
||||
if (typeof payload !== 'object' || payload === null) {
|
||||
return {}
|
||||
}
|
||||
const runtimeKind = 'runtimeKind' in payload ? payload.runtimeKind : undefined
|
||||
const runtimeVersion = 'runtimeVersion' in payload ? payload.runtimeVersion : undefined
|
||||
// Why drop other kinds: only a Node daemon reports here; anything else is an unknown peer.
|
||||
return {
|
||||
...(runtimeKind === 'node' ? { runtimeKind } : {}),
|
||||
...(typeof runtimeVersion === 'string' && runtimeVersion.length > 0 ? { runtimeVersion } : {})
|
||||
}
|
||||
}
|
||||
|
||||
function readPtySpawnHealthCoverage(
|
||||
payload: unknown,
|
||||
fallback: DaemonHealthCheck['coverage']
|
||||
): DaemonHealthCheck['coverage'] {
|
||||
if (typeof payload !== 'object' || payload === null) {
|
||||
return fallback
|
||||
}
|
||||
const coverage = 'coverage' in payload ? payload.coverage : undefined
|
||||
return coverage === 'pty-spawn' || coverage === 'handshake' ? coverage : fallback
|
||||
}
|
||||
|
||||
export function checkDaemonHealthWithCoverage(
|
||||
socketPath: string,
|
||||
tokenPath: string
|
||||
): Promise<DaemonHealthCheck> {
|
||||
return new Promise((resolve) => {
|
||||
// Older Windows daemons answered this RPC without spawning; an absent optional coverage
|
||||
// field must preserve that weaker meaning during adoption.
|
||||
const fallbackCoverage = process.platform === 'win32' ? 'handshake' : 'pty-spawn'
|
||||
const resolveVerdict = (verdict: DaemonHealth): void =>
|
||||
resolve({ verdict, coverage: fallbackCoverage })
|
||||
if (process.env[E2E_FORCE_DAEMON_HEALTH_UNREACHABLE_ENV] === '1') {
|
||||
resolve('unreachable')
|
||||
resolveVerdict('unreachable')
|
||||
return
|
||||
}
|
||||
|
||||
if (process.platform !== 'win32' && !existsSync(socketPath)) {
|
||||
resolve('unreachable')
|
||||
resolveVerdict('unreachable')
|
||||
return
|
||||
}
|
||||
|
||||
@@ -37,13 +79,13 @@ export function checkDaemonHealth(socketPath: string, tokenPath: string): Promis
|
||||
try {
|
||||
token = readFileSync(tokenPath, 'utf8').trim()
|
||||
} catch {
|
||||
resolve('unreachable')
|
||||
resolveVerdict('unreachable')
|
||||
return
|
||||
}
|
||||
|
||||
let settled = false
|
||||
let sock: Socket | null = null
|
||||
const settle = (result: DaemonHealth): void => {
|
||||
const settle = (result: DaemonHealthCheck): void => {
|
||||
if (settled) {
|
||||
return
|
||||
}
|
||||
@@ -58,7 +100,7 @@ export function checkDaemonHealth(socketPath: string, tokenPath: string): Promis
|
||||
sock?.off('connect', onConnect)
|
||||
sock?.off('data', onData)
|
||||
}
|
||||
const onError = (): void => settle('unreachable')
|
||||
const onError = (): void => settle({ verdict: 'unreachable', coverage: fallbackCoverage })
|
||||
const onConnect = (): void => {
|
||||
const hello: HelloMessage = {
|
||||
type: 'hello',
|
||||
@@ -89,13 +131,13 @@ export function checkDaemonHealth(socketPath: string, tokenPath: string): Promis
|
||||
try {
|
||||
message = JSON.parse(line) as Record<string, unknown>
|
||||
} catch {
|
||||
settle('rejected')
|
||||
settle({ verdict: 'rejected', coverage: fallbackCoverage })
|
||||
return
|
||||
}
|
||||
|
||||
if (message.type === 'hello') {
|
||||
if (!(message as HelloResponse).ok) {
|
||||
settle('rejected')
|
||||
settle({ verdict: 'rejected', coverage: fallbackCoverage })
|
||||
return
|
||||
}
|
||||
// Why: a protocol-live daemon with a stale cwd or node-pty helper
|
||||
@@ -106,12 +148,20 @@ export function checkDaemonHealth(socketPath: string, tokenPath: string): Promis
|
||||
}
|
||||
|
||||
if (message.id === 'health-1') {
|
||||
settle(message.ok === true ? 'healthy' : 'pty-spawn-unhealthy')
|
||||
const identity = readRuntimeIdentity(message.payload)
|
||||
settle({
|
||||
verdict: message.ok === true ? 'healthy' : 'pty-spawn-unhealthy',
|
||||
coverage: readPtySpawnHealthCoverage(message.payload, fallbackCoverage),
|
||||
...identity
|
||||
})
|
||||
return
|
||||
}
|
||||
}
|
||||
}
|
||||
const timer = setTimeout(() => settle('unreachable'), HEALTH_CHECK_TIMEOUT_MS)
|
||||
const timer = setTimeout(
|
||||
() => settle({ verdict: 'unreachable', coverage: fallbackCoverage }),
|
||||
HEALTH_CHECK_TIMEOUT_MS
|
||||
)
|
||||
|
||||
sock = connect({ path: socketPath })
|
||||
sock.on('error', onError)
|
||||
@@ -122,6 +172,13 @@ export function checkDaemonHealth(socketPath: string, tokenPath: string): Promis
|
||||
})
|
||||
}
|
||||
|
||||
export async function checkDaemonHealth(
|
||||
socketPath: string,
|
||||
tokenPath: string
|
||||
): Promise<DaemonHealth> {
|
||||
return (await checkDaemonHealthWithCoverage(socketPath, tokenPath)).verdict
|
||||
}
|
||||
|
||||
export async function healthCheckDaemon(socketPath: string, tokenPath: string): Promise<boolean> {
|
||||
return (await checkDaemonHealth(socketPath, tokenPath)) === 'healthy'
|
||||
}
|
||||
|
||||
@@ -442,6 +442,35 @@ describe('current daemon lifecycle retirement', () => {
|
||||
adopted.dispose()
|
||||
})
|
||||
|
||||
it('atomically retires an idle daemon and permanently fences adapter spawns', async () => {
|
||||
await startServer()
|
||||
const adapter = new DaemonPtyAdapter({ socketPath, tokenPath })
|
||||
|
||||
await expect(adapter.requestIdleRetirement()).resolves.toEqual({ state: 'retiring' })
|
||||
await expect(
|
||||
adapter.spawn({ sessionId: 'late-after-decommission', cols: 80, rows: 24 })
|
||||
).rejects.toThrow('Terminal daemon is decommissioning')
|
||||
await waitFor(() => onIdleShutdown.mock.calls.length === 1)
|
||||
adapter.dispose()
|
||||
})
|
||||
|
||||
it('reopens adapter admission when the daemon refuses retirement for a live session', async () => {
|
||||
await startServer()
|
||||
const adapter = new DaemonPtyAdapter({ socketPath, tokenPath })
|
||||
await adapter.spawn({ sessionId: 'already-live', cols: 80, rows: 24 })
|
||||
|
||||
await expect(adapter.requestIdleRetirement()).resolves.toEqual({
|
||||
state: 'busy',
|
||||
liveSessions: 1,
|
||||
admissionReopened: true
|
||||
})
|
||||
await expect(
|
||||
adapter.spawn({ sessionId: 'allowed-after-refusal', cols: 80, rows: 24 })
|
||||
).resolves.toMatchObject({ id: 'allowed-after-refusal' })
|
||||
expect(onIdleShutdown).not.toHaveBeenCalled()
|
||||
adapter.dispose()
|
||||
})
|
||||
|
||||
it('does not let repeated authenticated control probes extend the startup deadline', async () => {
|
||||
await startServer()
|
||||
const healthControl = connect(socketPath)
|
||||
|
||||
@@ -8,6 +8,8 @@ export {
|
||||
getDaemonEndpointFacts,
|
||||
getDaemonProvider,
|
||||
listLiveDaemonPtyIds,
|
||||
listLiveDaemonSessions,
|
||||
requestIdleDaemonRetirement,
|
||||
readDaemonPidRecord,
|
||||
replaceDaemonProvider,
|
||||
shutdownDaemon,
|
||||
|
||||
@@ -0,0 +1,42 @@
|
||||
import { afterEach, expect, it, vi } from 'vitest'
|
||||
|
||||
vi.mock('../ipc/pty', () => ({ setLocalPtyProvider: vi.fn() }))
|
||||
|
||||
import { DaemonPtyRouter } from './daemon-pty-router'
|
||||
import { createAdapter } from './daemon-pty-router-test-fixture'
|
||||
import {
|
||||
disconnectDaemon,
|
||||
listLiveDaemonSessions,
|
||||
replaceDaemonProvider,
|
||||
requestIdleDaemonRetirement
|
||||
} from './daemon-provider-state'
|
||||
import { PROTOCOL_VERSION } from './types'
|
||||
|
||||
afterEach(async () => {
|
||||
await disconnectDaemon()
|
||||
})
|
||||
|
||||
it('reads a census without an installed daemon as unverifiable, never as empty', async () => {
|
||||
await expect(listLiveDaemonSessions()).resolves.toBeNull()
|
||||
await expect(requestIdleDaemonRetirement()).resolves.toEqual({ state: 'unverifiable' })
|
||||
})
|
||||
|
||||
it('reads a census with an unanswered generation as unverifiable', async () => {
|
||||
const current = createAdapter('current', ['live-1'], undefined, PROTOCOL_VERSION)
|
||||
const legacy = createAdapter('legacy', [], undefined, PROTOCOL_VERSION)
|
||||
vi.mocked(legacy.listSessions).mockRejectedValue(new Error('daemon unreachable'))
|
||||
replaceDaemonProvider(new DaemonPtyRouter({ current, legacy: [legacy] }))
|
||||
|
||||
await expect(listLiveDaemonSessions()).resolves.toBeNull()
|
||||
})
|
||||
|
||||
it('lists every generation when each one answers', async () => {
|
||||
const current = createAdapter('current', ['live-1'], undefined, PROTOCOL_VERSION)
|
||||
const legacy = createAdapter('legacy', ['live-2'], undefined, PROTOCOL_VERSION)
|
||||
replaceDaemonProvider(new DaemonPtyRouter({ current, legacy: [legacy] }))
|
||||
|
||||
await expect(listLiveDaemonSessions()).resolves.toEqual([
|
||||
{ sessionId: 'live-1', isAlive: true },
|
||||
{ sessionId: 'live-2', isAlive: true }
|
||||
])
|
||||
})
|
||||
@@ -11,12 +11,16 @@ import {
|
||||
getMacDaemonTccAttributionHealth,
|
||||
type MacDaemonTccAttributionHealth
|
||||
} from './daemon-tcc-attribution'
|
||||
import { PROTOCOL_VERSION } from './types'
|
||||
import { PROTOCOL_VERSION, type SessionInfo } from './types'
|
||||
import type { DaemonIdleRetirementResult } from './daemon-pty-runtime-state'
|
||||
|
||||
let spawner: DaemonSpawner | null = null
|
||||
let adapter: DaemonProvider | null = null
|
||||
|
||||
export function installDaemonProvider(newSpawner: DaemonSpawner, newAdapter: DaemonProvider): void {
|
||||
export function installDaemonProvider(
|
||||
newSpawner: DaemonSpawner | null,
|
||||
newAdapter: DaemonProvider
|
||||
): void {
|
||||
spawner = newSpawner
|
||||
replaceDaemonProvider(newAdapter)
|
||||
}
|
||||
@@ -35,7 +39,12 @@ export function getDaemonSpawner(): DaemonSpawner | null {
|
||||
* from that state would be advertising recovery for terminals that cannot be recovered.
|
||||
*/
|
||||
export function daemonOwnsFreshPersistentPtys(): boolean {
|
||||
return adapter !== null && !(adapter instanceof DegradedDaemonPtyProvider)
|
||||
if (!adapter || adapter instanceof DegradedDaemonPtyProvider) {
|
||||
return false
|
||||
}
|
||||
return adapter instanceof DaemonPtyRouter
|
||||
? !adapter.getAllAdapters().some((entry) => entry.recoveryOnly)
|
||||
: !adapter.recoveryOnly
|
||||
}
|
||||
|
||||
/** Endpoint coordinates of the daemon this process installed, for out-of-band health probes. */
|
||||
@@ -95,26 +104,6 @@ export async function getCurrentDaemonMacTccAttributionHealth(): Promise<MacDaem
|
||||
)
|
||||
}
|
||||
|
||||
/** Returns null unless every daemon generation supplied an authoritative inventory. */
|
||||
export async function listLiveDaemonPtyIds(): Promise<string[] | null> {
|
||||
if (!adapter) {
|
||||
return null
|
||||
}
|
||||
const adapters =
|
||||
adapter instanceof DaemonPtyRouter || adapter instanceof DegradedDaemonPtyProvider
|
||||
? adapter.getAllAdapters()
|
||||
: [adapter]
|
||||
const inventories = await Promise.allSettled(
|
||||
adapters.map((daemonAdapter) => daemonAdapter.listProcesses())
|
||||
)
|
||||
if (inventories.some((inventory) => inventory.status === 'rejected')) {
|
||||
return null
|
||||
}
|
||||
return inventories.flatMap((inventory) =>
|
||||
inventory.status === 'fulfilled' ? inventory.value.map((process) => process.id) : []
|
||||
)
|
||||
}
|
||||
|
||||
// Why: keep the module-level adapter and ipc/pty.ts's localProvider in sync so app-quit can't dispose a stale reference.
|
||||
export function replaceDaemonProvider(newAdapter: DaemonProvider): void {
|
||||
adapter = newAdapter
|
||||
@@ -135,3 +124,57 @@ export async function shutdownDaemon(): Promise<void> {
|
||||
await spawner?.shutdown()
|
||||
spawner = null
|
||||
}
|
||||
|
||||
/** Returns null unless every daemon generation supplied an authoritative inventory. */
|
||||
export async function listLiveDaemonPtyIds(): Promise<string[] | null> {
|
||||
if (!adapter) {
|
||||
return null
|
||||
}
|
||||
const adapters =
|
||||
adapter instanceof DaemonPtyRouter || adapter instanceof DegradedDaemonPtyProvider
|
||||
? adapter.getAllAdapters()
|
||||
: [adapter]
|
||||
const inventories = await Promise.allSettled(
|
||||
adapters.map((daemonAdapter) => daemonAdapter.listProcesses())
|
||||
)
|
||||
if (inventories.some((inventory) => inventory.status === 'rejected')) {
|
||||
return null
|
||||
}
|
||||
return inventories.flatMap((inventory) =>
|
||||
inventory.status === 'fulfilled' ? inventory.value.map((process) => process.id) : []
|
||||
)
|
||||
}
|
||||
|
||||
/** Returns null unless every daemon generation supplied an authoritative session inventory. */
|
||||
export async function listLiveDaemonSessions(): Promise<SessionInfo[] | null> {
|
||||
if (!adapter) {
|
||||
return null
|
||||
}
|
||||
const adapters =
|
||||
adapter instanceof DaemonPtyRouter || adapter instanceof DegradedDaemonPtyProvider
|
||||
? adapter.getAllAdapters()
|
||||
: [adapter]
|
||||
const inventories = await Promise.allSettled(
|
||||
adapters.map((daemonAdapter) => daemonAdapter.listSessions())
|
||||
)
|
||||
if (inventories.some((inventory) => inventory.status === 'rejected')) {
|
||||
return null
|
||||
}
|
||||
return inventories.flatMap((inventory) =>
|
||||
inventory.status === 'fulfilled' ? inventory.value : []
|
||||
)
|
||||
}
|
||||
|
||||
/** Atomically fence new daemon terminals and retire only an idle, single-generation daemon. */
|
||||
export async function requestIdleDaemonRetirement(): Promise<DaemonIdleRetirementResult> {
|
||||
if (!adapter) {
|
||||
return { state: 'unverifiable' }
|
||||
}
|
||||
if (adapter instanceof DegradedDaemonPtyProvider) {
|
||||
return { state: 'unverifiable' }
|
||||
}
|
||||
if (adapter instanceof DaemonPtyRouter) {
|
||||
return adapter.requestIdleRetirement()
|
||||
}
|
||||
return adapter.requestIdleRetirement()
|
||||
}
|
||||
|
||||
@@ -3,7 +3,12 @@ import { removeDaemonListener } from './daemon-listener-registry'
|
||||
import { emitPtyListeners } from './daemon-pty-listener-emission'
|
||||
import type { PtyIncarnationId } from '../../shared/pty-incarnation'
|
||||
import { DaemonPtySessionInventory } from './daemon-pty-session-inventory'
|
||||
import { CLEAN_DISCONNECT_PROTOCOL_VERSION } from './types'
|
||||
import {
|
||||
CLEAN_DISCONNECT_PROTOCOL_VERSION,
|
||||
type ListSessionsResult,
|
||||
type ShutdownIfIdleResult
|
||||
} from './types'
|
||||
import type { DaemonIdleRetirementResult } from './daemon-pty-runtime-state'
|
||||
import type { PtyBackgroundStreamEvent } from '../providers/types'
|
||||
|
||||
export abstract class DaemonPtyEventSubscriptions extends DaemonPtySessionInventory {
|
||||
@@ -85,6 +90,68 @@ export abstract class DaemonPtyEventSubscriptions extends DaemonPtySessionInvent
|
||||
this.recordAuthenticatedIdentity()
|
||||
}
|
||||
|
||||
async requestIdleRetirement(): Promise<DaemonIdleRetirementResult> {
|
||||
if (this.protocolVersion < CLEAN_DISCONNECT_PROTOCOL_VERSION) {
|
||||
return { state: 'unsupported' }
|
||||
}
|
||||
if (this.idleRetirementState === 'retiring') {
|
||||
return { state: 'retiring' }
|
||||
}
|
||||
if (this.idleRetirementPromise) {
|
||||
return this.idleRetirementPromise
|
||||
}
|
||||
if (
|
||||
this.disconnectOnlyPromise ||
|
||||
(this.respawnAdoptionClosed && this.idleRetirementState === 'open')
|
||||
) {
|
||||
return { state: 'unverifiable' }
|
||||
}
|
||||
this.idleRetirementAdmissionClosed = true
|
||||
this.respawnAdoptionClosed = true
|
||||
this.idleRetirementState = 'checking'
|
||||
const request = this.finishIdleRetirementRequest().finally(() => {
|
||||
if (this.idleRetirementPromise === request) {
|
||||
this.idleRetirementPromise = null
|
||||
}
|
||||
})
|
||||
this.idleRetirementPromise = request
|
||||
return request
|
||||
}
|
||||
|
||||
private async finishIdleRetirementRequest(): Promise<DaemonIdleRetirementResult> {
|
||||
try {
|
||||
await this.client.ensureConnected()
|
||||
const result = await this.client.request<ShutdownIfIdleResult>('shutdownIfIdle', undefined)
|
||||
if (result.retiring) {
|
||||
this.idleRetirementState = 'retiring'
|
||||
return { state: 'retiring' }
|
||||
}
|
||||
let liveSessions: number | null = null
|
||||
try {
|
||||
const inventory = await this.client.request<ListSessionsResult>('listSessions', undefined)
|
||||
liveSessions = inventory.sessions.filter((session) => session.isAlive).length
|
||||
} catch {
|
||||
liveSessions = null
|
||||
}
|
||||
this.reopenAfterRefusedIdleRetirement()
|
||||
return {
|
||||
state: 'busy',
|
||||
liveSessions,
|
||||
...(!this.recoveryOnly ? { admissionReopened: true as const } : {})
|
||||
}
|
||||
} catch {
|
||||
// The daemon may have accepted before contact was lost; keep admission and respawn fenced.
|
||||
this.idleRetirementState = 'unverifiable'
|
||||
return { state: 'unverifiable' }
|
||||
}
|
||||
}
|
||||
|
||||
private reopenAfterRefusedIdleRetirement(): void {
|
||||
this.idleRetirementState = 'open'
|
||||
this.idleRetirementAdmissionClosed = false
|
||||
this.respawnAdoptionClosed = false
|
||||
}
|
||||
|
||||
// Why: unlike dispose(), leave history files unclean (no endedAt) so the next launch treats them as crash-recoverable,
|
||||
// but still write a final checkpoint so a daemon crash while Orca is closed has recovery data.
|
||||
async disconnectOnly(): Promise<void> {
|
||||
|
||||
@@ -127,7 +127,7 @@ export abstract class DaemonPtyProcessInspection extends DaemonPtyBufferSnapshot
|
||||
// Why: an unminted session id (worktreeId === null) can't be tied to a live worktree, so it's treated as an orphan.
|
||||
const { worktreeId } = parsePtySessionId(session.sessionId)
|
||||
|
||||
if (worktreeId === null || !validWorktreeIds.has(worktreeId)) {
|
||||
if (!this.recoveryOnly && (worktreeId === null || !validWorktreeIds.has(worktreeId))) {
|
||||
try {
|
||||
await this.client.request('kill', { sessionId: session.sessionId })
|
||||
} catch {
|
||||
|
||||
@@ -0,0 +1,167 @@
|
||||
import { vi } from 'vitest'
|
||||
import { settledWriteStub } from '../providers/settled-pty-write-stub'
|
||||
import type { DaemonPtyAdapter } from './daemon-pty-adapter'
|
||||
import type { PtyBackgroundStreamEvent, PtySpawnOptions, PtySpawnResult } from '../providers/types'
|
||||
import {
|
||||
AGENT_SESSION_CLAIM_DAEMON_PROTOCOL_VERSION,
|
||||
AGENT_SESSION_CREATE_OPERATION_DAEMON_PROTOCOL_VERSION,
|
||||
GIT_CREDENTIAL_GUARD_HOST_PROTOCOL_VERSION
|
||||
} from './types'
|
||||
import { SNAPSHOT_SERIALIZER_FIDELITY_DAEMON_PROTOCOL_VERSION } from './daemon-protocol-version'
|
||||
|
||||
type AdapterMock = DaemonPtyAdapter & {
|
||||
emitData: (id: string, data: string, sequenceChars?: number) => void
|
||||
emitBackground: (event: PtyBackgroundStreamEvent) => void
|
||||
emitExit: (id: string, code: number, incarnationId?: string) => void
|
||||
emitIdentityChange: () => void
|
||||
triggerWriteUnavailable: (id: string) => void
|
||||
}
|
||||
|
||||
export function createAdapter(
|
||||
label: string,
|
||||
sessions: string[] = [],
|
||||
reconcileResult?: { alive: string[]; killed: string[] },
|
||||
protocolVersion = GIT_CREDENTIAL_GUARD_HOST_PROTOCOL_VERSION
|
||||
): AdapterMock {
|
||||
const writes: { id: string; data: string }[] = []
|
||||
const dataListeners: ((payload: { id: string; data: string; sequenceChars?: number }) => void)[] =
|
||||
[]
|
||||
const backgroundListeners: ((payload: PtyBackgroundStreamEvent) => void)[] = []
|
||||
const writeUnavailableListeners: ((payload: { id: string }) => void)[] = []
|
||||
const exitListeners: ((payload: { id: string; code: number; incarnationId?: string }) => void)[] =
|
||||
[]
|
||||
const identityChangeListeners: (() => void)[] = []
|
||||
// oxlint-disable-next-line typescript/consistent-type-assertions -- SAFETY: the router calls only the adapter members this mock defines.
|
||||
return {
|
||||
protocolVersion,
|
||||
supportsGitCredentialGuardHost: () =>
|
||||
protocolVersion >= GIT_CREDENTIAL_GUARD_HOST_PROTOCOL_VERSION,
|
||||
supportsAgentSessionClaims: () =>
|
||||
protocolVersion >= AGENT_SESSION_CLAIM_DAEMON_PROTOCOL_VERSION,
|
||||
supportsAgentSessionCreateOperations: () =>
|
||||
protocolVersion >= AGENT_SESSION_CREATE_OPERATION_DAEMON_PROTOCOL_VERSION,
|
||||
providesAgentSessionOwnerListings: () =>
|
||||
protocolVersion >= AGENT_SESSION_CLAIM_DAEMON_PROTOCOL_VERSION,
|
||||
canProvideAuthoritativeBufferSnapshot: () =>
|
||||
protocolVersion >= SNAPSHOT_SERIALIZER_FIDELITY_DAEMON_PROTOCOL_VERSION,
|
||||
spawn: vi.fn(async (opts: PtySpawnOptions): Promise<PtySpawnResult> => {
|
||||
const id = opts.sessionId ?? `${label}-new`
|
||||
sessions.push(id)
|
||||
return { id }
|
||||
}),
|
||||
listProcesses: vi.fn(async () =>
|
||||
sessions.map((id) => ({
|
||||
id,
|
||||
cwd: '',
|
||||
title: label
|
||||
}))
|
||||
),
|
||||
listSessions: vi.fn(async () => sessions.map((sessionId) => ({ sessionId, isAlive: true }))),
|
||||
requestIdleRetirement: vi.fn(async () => ({ state: 'retiring' as const })),
|
||||
hasPty: vi.fn((id: string) => sessions.includes(id)),
|
||||
probePtyLiveness: vi.fn(async (id: string) => sessions.includes(id)),
|
||||
write: vi.fn((id: string, data: string) => {
|
||||
writes.push({ id, data })
|
||||
}),
|
||||
writeWithSettlement: vi.fn(settledWriteStub()),
|
||||
resize: vi.fn(),
|
||||
setPtyBackgrounded: vi.fn(),
|
||||
getBufferSnapshot: vi.fn(async () => null),
|
||||
shutdown: vi.fn(async (id: string) => {
|
||||
const idx = sessions.indexOf(id)
|
||||
if (idx !== -1) {
|
||||
sessions.splice(idx, 1)
|
||||
}
|
||||
}),
|
||||
attach: vi.fn(async () => {}),
|
||||
sendSignal: vi.fn(async () => {}),
|
||||
getCwd: vi.fn(async () => ''),
|
||||
getInitialCwd: vi.fn(async () => ''),
|
||||
clearBuffer: vi.fn(async () => {}),
|
||||
acknowledgeDataEvent: vi.fn(),
|
||||
hasChildProcesses: vi.fn(async () => false),
|
||||
getForegroundProcess: vi.fn(async () => null),
|
||||
inspectProcess: vi.fn(async () => ({ foregroundProcess: null, hasChildProcesses: false })),
|
||||
confirmForegroundProcess: vi.fn(async () => `${label}-confirmed`),
|
||||
serialize: vi.fn(async () => '{}'),
|
||||
revive: vi.fn(async () => {}),
|
||||
getDefaultShell: vi.fn(async () => '/bin/zsh'),
|
||||
getProfiles: vi.fn(async () => []),
|
||||
onData: vi.fn(
|
||||
(callback: (payload: { id: string; data: string; sequenceChars?: number }) => void) => {
|
||||
dataListeners.push(callback)
|
||||
return () => {
|
||||
const idx = dataListeners.indexOf(callback)
|
||||
if (idx !== -1) {
|
||||
dataListeners.splice(idx, 1)
|
||||
}
|
||||
}
|
||||
}
|
||||
),
|
||||
onBackgroundStreamEvent: vi.fn((callback: (payload: PtyBackgroundStreamEvent) => void) => {
|
||||
backgroundListeners.push(callback)
|
||||
return () => {
|
||||
const idx = backgroundListeners.indexOf(callback)
|
||||
if (idx !== -1) {
|
||||
backgroundListeners.splice(idx, 1)
|
||||
}
|
||||
}
|
||||
}),
|
||||
onWriteUnavailable: vi.fn((callback: (payload: { id: string }) => void) => {
|
||||
writeUnavailableListeners.push(callback)
|
||||
return () => {
|
||||
const idx = writeUnavailableListeners.indexOf(callback)
|
||||
if (idx !== -1) {
|
||||
writeUnavailableListeners.splice(idx, 1)
|
||||
}
|
||||
}
|
||||
}),
|
||||
onExit: vi.fn(
|
||||
(callback: (payload: { id: string; code: number; incarnationId?: string }) => void) => {
|
||||
exitListeners.push(callback)
|
||||
return () => {
|
||||
const idx = exitListeners.indexOf(callback)
|
||||
if (idx !== -1) {
|
||||
exitListeners.splice(idx, 1)
|
||||
}
|
||||
}
|
||||
}
|
||||
),
|
||||
onDaemonIdentityChanged: vi.fn((callback: () => void) => {
|
||||
identityChangeListeners.push(callback)
|
||||
return () => {
|
||||
const idx = identityChangeListeners.indexOf(callback)
|
||||
if (idx !== -1) {
|
||||
identityChangeListeners.splice(idx, 1)
|
||||
}
|
||||
}
|
||||
}),
|
||||
ackColdRestore: vi.fn(),
|
||||
clearTombstone: vi.fn(),
|
||||
reconcileOnStartup: vi.fn(async () => reconcileResult ?? { alive: sessions, killed: [] }),
|
||||
dispose: vi.fn(),
|
||||
disconnectOnly: vi.fn(async () => {}),
|
||||
emitData: (id: string, data: string, sequenceChars?: number) => {
|
||||
for (const listener of dataListeners) {
|
||||
listener({ id, data, ...(sequenceChars === undefined ? {} : { sequenceChars }) })
|
||||
}
|
||||
},
|
||||
emitBackground: (event: PtyBackgroundStreamEvent) => {
|
||||
for (const listener of backgroundListeners) {
|
||||
listener(event)
|
||||
}
|
||||
},
|
||||
emitExit: (id: string, code: number, incarnationId?: string) => {
|
||||
for (const listener of exitListeners) {
|
||||
listener({ id, code, ...(incarnationId ? { incarnationId } : {}) })
|
||||
}
|
||||
},
|
||||
emitIdentityChange: () => identityChangeListeners.forEach((listener) => listener()),
|
||||
triggerWriteUnavailable: (id: string) => {
|
||||
for (const listener of writeUnavailableListeners) {
|
||||
listener({ id })
|
||||
}
|
||||
},
|
||||
_writes: writes
|
||||
} as unknown as AdapterMock
|
||||
}
|
||||
@@ -1,29 +1,20 @@
|
||||
import { createAdapter } from './daemon-pty-router-test-fixture'
|
||||
import { describe, expect, it, vi } from 'vitest'
|
||||
import { DaemonPtyRouter } from './daemon-pty-router'
|
||||
import { stubWriteSettlement } from '../providers/settled-pty-write-stub'
|
||||
import { SessionNotFoundError, TerminalSessionOwnerUnverifiedError } from './daemon-errors'
|
||||
import type { DaemonPtyAdapter } from './daemon-pty-adapter'
|
||||
import { settledWriteStub, stubWriteSettlement } from '../providers/settled-pty-write-stub'
|
||||
import type { PtyBackgroundStreamEvent, PtySpawnOptions, PtySpawnResult } from '../providers/types'
|
||||
import type { PtySpawnResult } from '../providers/types'
|
||||
import {
|
||||
AGENT_SESSION_CLAIM_DAEMON_PROTOCOL_VERSION,
|
||||
AGENT_SESSION_CREATE_OPERATION_DAEMON_PROTOCOL_VERSION,
|
||||
GIT_CREDENTIAL_GUARD_HOST_PROTOCOL_VERSION
|
||||
AGENT_SESSION_CREATE_OPERATION_DAEMON_PROTOCOL_VERSION
|
||||
} from './types'
|
||||
import {
|
||||
HISTORY_SEED_TRANSFER_PROTOCOL_VERSION,
|
||||
PROTOCOL_VERSION,
|
||||
SNAPSHOT_SERIALIZER_FIDELITY_DAEMON_PROTOCOL_VERSION,
|
||||
STABLE_PANE_ATTACH_ONLY_DAEMON_PROTOCOL_VERSION
|
||||
} from './daemon-protocol-version'
|
||||
|
||||
type AdapterMock = DaemonPtyAdapter & {
|
||||
emitData: (id: string, data: string, sequenceChars?: number) => void
|
||||
emitBackground: (event: PtyBackgroundStreamEvent) => void
|
||||
emitExit: (id: string, code: number, incarnationId?: string) => void
|
||||
emitIdentityChange: () => void
|
||||
triggerWriteUnavailable: (id: string) => void
|
||||
}
|
||||
|
||||
const LARGE_RECONCILE_SESSION_COUNT = 150_000
|
||||
|
||||
function buildSessionIds(prefix: string, count: number): string[] {
|
||||
@@ -34,152 +25,6 @@ function buildSessionIds(prefix: string, count: number): string[] {
|
||||
return ids
|
||||
}
|
||||
|
||||
function createAdapter(
|
||||
label: string,
|
||||
sessions: string[] = [],
|
||||
reconcileResult?: { alive: string[]; killed: string[] },
|
||||
protocolVersion = GIT_CREDENTIAL_GUARD_HOST_PROTOCOL_VERSION
|
||||
): AdapterMock {
|
||||
const writes: { id: string; data: string }[] = []
|
||||
const dataListeners: ((payload: { id: string; data: string; sequenceChars?: number }) => void)[] =
|
||||
[]
|
||||
const backgroundListeners: ((payload: PtyBackgroundStreamEvent) => void)[] = []
|
||||
const writeUnavailableListeners: ((payload: { id: string }) => void)[] = []
|
||||
const exitListeners: ((payload: { id: string; code: number; incarnationId?: string }) => void)[] =
|
||||
[]
|
||||
const identityChangeListeners: (() => void)[] = []
|
||||
return {
|
||||
protocolVersion,
|
||||
supportsGitCredentialGuardHost: () =>
|
||||
protocolVersion >= GIT_CREDENTIAL_GUARD_HOST_PROTOCOL_VERSION,
|
||||
supportsAgentSessionClaims: () =>
|
||||
protocolVersion >= AGENT_SESSION_CLAIM_DAEMON_PROTOCOL_VERSION,
|
||||
supportsAgentSessionCreateOperations: () =>
|
||||
protocolVersion >= AGENT_SESSION_CREATE_OPERATION_DAEMON_PROTOCOL_VERSION,
|
||||
providesAgentSessionOwnerListings: () =>
|
||||
protocolVersion >= AGENT_SESSION_CLAIM_DAEMON_PROTOCOL_VERSION,
|
||||
canProvideAuthoritativeBufferSnapshot: () =>
|
||||
protocolVersion >= SNAPSHOT_SERIALIZER_FIDELITY_DAEMON_PROTOCOL_VERSION,
|
||||
spawn: vi.fn(async (opts: PtySpawnOptions): Promise<PtySpawnResult> => {
|
||||
const id = opts.sessionId ?? `${label}-new`
|
||||
sessions.push(id)
|
||||
return { id }
|
||||
}),
|
||||
listProcesses: vi.fn(async () =>
|
||||
sessions.map((id) => ({
|
||||
id,
|
||||
cwd: '',
|
||||
title: label
|
||||
}))
|
||||
),
|
||||
hasPty: vi.fn((id: string) => sessions.includes(id)),
|
||||
probePtyLiveness: vi.fn(async (id: string) => sessions.includes(id)),
|
||||
write: vi.fn((id: string, data: string) => {
|
||||
writes.push({ id, data })
|
||||
}),
|
||||
writeWithSettlement: vi.fn(settledWriteStub()),
|
||||
resize: vi.fn(),
|
||||
setPtyBackgrounded: vi.fn(),
|
||||
getBufferSnapshot: vi.fn(async () => null),
|
||||
shutdown: vi.fn(async (id: string) => {
|
||||
const idx = sessions.indexOf(id)
|
||||
if (idx !== -1) {
|
||||
sessions.splice(idx, 1)
|
||||
}
|
||||
}),
|
||||
attach: vi.fn(async () => {}),
|
||||
sendSignal: vi.fn(async () => {}),
|
||||
getCwd: vi.fn(async () => ''),
|
||||
getInitialCwd: vi.fn(async () => ''),
|
||||
clearBuffer: vi.fn(async () => {}),
|
||||
acknowledgeDataEvent: vi.fn(),
|
||||
hasChildProcesses: vi.fn(async () => false),
|
||||
getForegroundProcess: vi.fn(async () => null),
|
||||
inspectProcess: vi.fn(async () => ({ foregroundProcess: null, hasChildProcesses: false })),
|
||||
confirmForegroundProcess: vi.fn(async () => `${label}-confirmed`),
|
||||
serialize: vi.fn(async () => '{}'),
|
||||
revive: vi.fn(async () => {}),
|
||||
getDefaultShell: vi.fn(async () => '/bin/zsh'),
|
||||
getProfiles: vi.fn(async () => []),
|
||||
onData: vi.fn(
|
||||
(callback: (payload: { id: string; data: string; sequenceChars?: number }) => void) => {
|
||||
dataListeners.push(callback)
|
||||
return () => {
|
||||
const idx = dataListeners.indexOf(callback)
|
||||
if (idx !== -1) {
|
||||
dataListeners.splice(idx, 1)
|
||||
}
|
||||
}
|
||||
}
|
||||
),
|
||||
onBackgroundStreamEvent: vi.fn((callback: (payload: PtyBackgroundStreamEvent) => void) => {
|
||||
backgroundListeners.push(callback)
|
||||
return () => {
|
||||
const idx = backgroundListeners.indexOf(callback)
|
||||
if (idx !== -1) {
|
||||
backgroundListeners.splice(idx, 1)
|
||||
}
|
||||
}
|
||||
}),
|
||||
onWriteUnavailable: vi.fn((callback: (payload: { id: string }) => void) => {
|
||||
writeUnavailableListeners.push(callback)
|
||||
return () => {
|
||||
const idx = writeUnavailableListeners.indexOf(callback)
|
||||
if (idx !== -1) {
|
||||
writeUnavailableListeners.splice(idx, 1)
|
||||
}
|
||||
}
|
||||
}),
|
||||
onExit: vi.fn(
|
||||
(callback: (payload: { id: string; code: number; incarnationId?: string }) => void) => {
|
||||
exitListeners.push(callback)
|
||||
return () => {
|
||||
const idx = exitListeners.indexOf(callback)
|
||||
if (idx !== -1) {
|
||||
exitListeners.splice(idx, 1)
|
||||
}
|
||||
}
|
||||
}
|
||||
),
|
||||
onDaemonIdentityChanged: vi.fn((callback: () => void) => {
|
||||
identityChangeListeners.push(callback)
|
||||
return () => {
|
||||
const idx = identityChangeListeners.indexOf(callback)
|
||||
if (idx !== -1) {
|
||||
identityChangeListeners.splice(idx, 1)
|
||||
}
|
||||
}
|
||||
}),
|
||||
ackColdRestore: vi.fn(),
|
||||
clearTombstone: vi.fn(),
|
||||
reconcileOnStartup: vi.fn(async () => reconcileResult ?? { alive: sessions, killed: [] }),
|
||||
dispose: vi.fn(),
|
||||
disconnectOnly: vi.fn(async () => {}),
|
||||
emitData: (id: string, data: string, sequenceChars?: number) => {
|
||||
for (const listener of dataListeners) {
|
||||
listener({ id, data, ...(sequenceChars === undefined ? {} : { sequenceChars }) })
|
||||
}
|
||||
},
|
||||
emitBackground: (event: PtyBackgroundStreamEvent) => {
|
||||
for (const listener of backgroundListeners) {
|
||||
listener(event)
|
||||
}
|
||||
},
|
||||
emitExit: (id: string, code: number, incarnationId?: string) => {
|
||||
for (const listener of exitListeners) {
|
||||
listener({ id, code, ...(incarnationId ? { incarnationId } : {}) })
|
||||
}
|
||||
},
|
||||
emitIdentityChange: () => identityChangeListeners.forEach((listener) => listener()),
|
||||
triggerWriteUnavailable: (id: string) => {
|
||||
for (const listener of writeUnavailableListeners) {
|
||||
listener({ id })
|
||||
}
|
||||
},
|
||||
_writes: writes
|
||||
} as unknown as AdapterMock
|
||||
}
|
||||
|
||||
it('forwards dead-endpoint write-unavailable signals from every routed adapter', () => {
|
||||
// Why revert-sensitive: main subscribes on the ROUTED provider, so if the router
|
||||
// does not forward this the STA-2373 fan-out never reaches the renderer and only
|
||||
@@ -247,6 +92,83 @@ it('forwards the owning legacy daemon sequence from attach', async () => {
|
||||
})
|
||||
|
||||
describe('DaemonPtyRouter', () => {
|
||||
describe('idle retirement', () => {
|
||||
it('retires every empty daemon generation and fences subsequent spawns', async () => {
|
||||
const current = createAdapter('current', [], undefined, PROTOCOL_VERSION)
|
||||
const legacy = createAdapter('legacy', [], undefined, PROTOCOL_VERSION)
|
||||
const router = new DaemonPtyRouter({ current, legacy: [legacy] })
|
||||
|
||||
await expect(router.requestIdleRetirement()).resolves.toEqual({ state: 'retiring' })
|
||||
expect(current.requestIdleRetirement).toHaveBeenCalledOnce()
|
||||
expect(legacy.requestIdleRetirement).toHaveBeenCalledOnce()
|
||||
await expect(router.spawn({ sessionId: 'late', cols: 80, rows: 24 })).rejects.toThrow(
|
||||
'Terminal daemon is decommissioning'
|
||||
)
|
||||
})
|
||||
|
||||
it('reports live inventory before retiring any generation and reopens admission', async () => {
|
||||
const current = createAdapter('current', [], undefined, PROTOCOL_VERSION)
|
||||
const legacy = createAdapter('legacy', ['legacy-live'], undefined, PROTOCOL_VERSION)
|
||||
const router = new DaemonPtyRouter({ current, legacy: [legacy] })
|
||||
|
||||
await expect(router.requestIdleRetirement()).resolves.toEqual({
|
||||
state: 'busy',
|
||||
liveSessions: 1,
|
||||
admissionReopened: true
|
||||
})
|
||||
expect(current.requestIdleRetirement).not.toHaveBeenCalled()
|
||||
expect(legacy.requestIdleRetirement).not.toHaveBeenCalled()
|
||||
await expect(
|
||||
router.spawn({ sessionId: 'after-refusal', cols: 80, rows: 24 })
|
||||
).resolves.toEqual({
|
||||
id: 'after-refusal'
|
||||
})
|
||||
})
|
||||
|
||||
it('does not partially retire when a generation predates clean idle shutdown', async () => {
|
||||
const current = createAdapter('current', [], undefined, PROTOCOL_VERSION)
|
||||
const legacy = createAdapter('legacy', [], undefined, 23)
|
||||
const router = new DaemonPtyRouter({ current, legacy: [legacy] })
|
||||
|
||||
await expect(router.requestIdleRetirement()).resolves.toEqual({ state: 'unsupported' })
|
||||
expect(current.requestIdleRetirement).not.toHaveBeenCalled()
|
||||
expect(legacy.requestIdleRetirement).not.toHaveBeenCalled()
|
||||
})
|
||||
|
||||
it('keeps admission fenced after a partial multi-generation retirement', async () => {
|
||||
const current = createAdapter('current', [], undefined, PROTOCOL_VERSION)
|
||||
const legacy = createAdapter('legacy', [], undefined, PROTOCOL_VERSION)
|
||||
vi.mocked(legacy.requestIdleRetirement).mockResolvedValueOnce({
|
||||
state: 'busy',
|
||||
liveSessions: 0
|
||||
})
|
||||
const router = new DaemonPtyRouter({ current, legacy: [legacy] })
|
||||
|
||||
await expect(router.requestIdleRetirement()).resolves.toEqual({ state: 'unverifiable' })
|
||||
await expect(router.spawn({ sessionId: 'unsafe', cols: 80, rows: 24 })).rejects.toThrow(
|
||||
'Terminal daemon is decommissioning'
|
||||
)
|
||||
})
|
||||
|
||||
it('does not certify reopened admission when another generation retired beside live sessions', async () => {
|
||||
const current = createAdapter('current', [], undefined, PROTOCOL_VERSION)
|
||||
const legacy = createAdapter('legacy', [], undefined, PROTOCOL_VERSION)
|
||||
vi.mocked(legacy.requestIdleRetirement).mockResolvedValueOnce({
|
||||
state: 'busy',
|
||||
liveSessions: 1,
|
||||
admissionReopened: true
|
||||
})
|
||||
const router = new DaemonPtyRouter({ current, legacy: [legacy] })
|
||||
await expect(router.requestIdleRetirement()).resolves.toEqual({
|
||||
state: 'busy',
|
||||
liveSessions: 1
|
||||
})
|
||||
await expect(
|
||||
router.spawn({ sessionId: 'unsafe-partial', cols: 80, rows: 24 })
|
||||
).rejects.toThrow('Terminal daemon is decommissioning')
|
||||
})
|
||||
})
|
||||
|
||||
it('reports separate conservative resume and fresh-create boundaries', () => {
|
||||
const current = createAdapter(
|
||||
'current',
|
||||
|
||||
@@ -1,3 +1,4 @@
|
||||
import { reconcileDaemonRouterSessions } from './daemon-router-session-reconciliation'
|
||||
import type { DaemonPtyAdapter } from './daemon-pty-adapter'
|
||||
import { DaemonPtyAdapterSubscriptionFanout } from './daemon-pty-adapter-subscription-fanout'
|
||||
import type {
|
||||
@@ -12,6 +13,8 @@ import type { PtyProcessInspection } from '../providers/pty-process-inspection'
|
||||
import { shouldHandoffDaemonHistory } from './daemon-history-handoff'
|
||||
import type { DaemonPtyRouterDataEvent, DaemonPtyRouterExitEvent } from './daemon-pty-router-events'
|
||||
import { DaemonSessionOwnerResolver } from './daemon-session-owner-resolution'
|
||||
import type { DaemonIdleRetirementResult } from './daemon-pty-runtime-state'
|
||||
import { DaemonRouterRetirement } from './daemon-router-retirement'
|
||||
import type { WriteSettlement } from '../../shared/pty-write-settlement'
|
||||
import type { TerminalOscColorQueryReplyColors } from '../../shared/terminal-osc-color-reply'
|
||||
|
||||
@@ -21,6 +24,7 @@ export class DaemonPtyRouter implements IPtyProvider {
|
||||
private sessionAdapters = new Map<string, DaemonPtyAdapter>()
|
||||
private readonly ownerResolver: DaemonSessionOwnerResolver<DaemonPtyAdapter>
|
||||
private readonly subscriptions: DaemonPtyAdapterSubscriptionFanout
|
||||
private readonly retirement = new DaemonRouterRetirement(() => this.allAdapters())
|
||||
|
||||
constructor(opts: { current: DaemonPtyAdapter; legacy: DaemonPtyAdapter[] }) {
|
||||
this.current = opts.current
|
||||
@@ -40,17 +44,30 @@ export class DaemonPtyRouter implements IPtyProvider {
|
||||
}
|
||||
|
||||
async spawn(opts: PtySpawnOptions): Promise<PtySpawnResult> {
|
||||
if (opts.attachOnly && opts.sessionId) {
|
||||
return await this.ownerResolver.spawnAttachOnly({ ...opts, sessionId: opts.sessionId })
|
||||
if (this.retirement.admissionClosed) {
|
||||
throw new Error('Terminal daemon is decommissioning')
|
||||
}
|
||||
const adapter = opts.sessionId ? this.sessionAdapters.get(opts.sessionId) : undefined
|
||||
const target = adapter ?? this.current
|
||||
const result = await target.spawn(opts)
|
||||
// Why: the adapter filters intentional recovery exits and canonical-ID races before publishing proof.
|
||||
if (!result.exitedBeforeSpawnReply) {
|
||||
this.ownerResolver.recordRoute(result.id, target, result.incarnationId)
|
||||
// Why counted: an idle-retirement census must not race a spawn it cannot yet see.
|
||||
this.retirement.spawnInFlight++
|
||||
try {
|
||||
if (opts.attachOnly && opts.sessionId) {
|
||||
return await this.ownerResolver.spawnAttachOnly({ ...opts, sessionId: opts.sessionId })
|
||||
}
|
||||
const adapter = opts.sessionId ? this.sessionAdapters.get(opts.sessionId) : undefined
|
||||
const target = adapter ?? this.current
|
||||
const result = await target.spawn(opts)
|
||||
// Why: the adapter filters intentional recovery exits and canonical-ID races before publishing proof.
|
||||
if (!result.exitedBeforeSpawnReply) {
|
||||
this.ownerResolver.recordRoute(result.id, target, result.incarnationId)
|
||||
}
|
||||
return result
|
||||
} finally {
|
||||
this.retirement.spawnInFlight--
|
||||
}
|
||||
return result
|
||||
}
|
||||
|
||||
requestIdleRetirement(): Promise<DaemonIdleRetirementResult> {
|
||||
return this.retirement.requestIdleRetirement()
|
||||
}
|
||||
|
||||
supportsGitCredentialGuardHost(sessionId?: string): boolean {
|
||||
@@ -257,38 +274,10 @@ export class DaemonPtyRouter implements IPtyProvider {
|
||||
this.adapterFor(sessionId).clearTombstone(sessionId)
|
||||
}
|
||||
|
||||
async reconcileOnStartup(validWorktreeIds: Set<string>): Promise<{
|
||||
alive: string[]
|
||||
killed: string[]
|
||||
}> {
|
||||
const alive: string[] = []
|
||||
const killed: string[] = []
|
||||
const aliveProviders = new Map<string, Set<DaemonPtyAdapter>>()
|
||||
for (const adapter of this.allAdapters()) {
|
||||
const result = await adapter.reconcileOnStartup(validWorktreeIds)
|
||||
// Why: daemon startup can reconcile many restored sessions; spreading
|
||||
// those arrays into push can exceed JavaScript's argument limit.
|
||||
for (const id of result.alive) {
|
||||
alive.push(id)
|
||||
}
|
||||
for (const id of result.killed) {
|
||||
killed.push(id)
|
||||
}
|
||||
for (const id of result.alive) {
|
||||
const providers = aliveProviders.get(id) ?? new Set<DaemonPtyAdapter>()
|
||||
providers.add(adapter)
|
||||
aliveProviders.set(id, providers)
|
||||
}
|
||||
}
|
||||
for (const id of new Set([...alive, ...killed])) {
|
||||
const providers = aliveProviders.get(id)
|
||||
if (providers?.size === 1) {
|
||||
this.ownerResolver.recordRoute(id, providers.values().next().value!)
|
||||
} else {
|
||||
this.ownerResolver.forgetRoute(id)
|
||||
}
|
||||
}
|
||||
return { alive, killed }
|
||||
async reconcileOnStartup(
|
||||
validWorktreeIds: Set<string>
|
||||
): Promise<{ alive: string[]; killed: string[] }> {
|
||||
return reconcileDaemonRouterSessions(this.allAdapters(), this.ownerResolver, validWorktreeIds)
|
||||
}
|
||||
|
||||
dispose(): void {
|
||||
|
||||
@@ -57,6 +57,7 @@ export type DaemonPtyAdapterOptions = {
|
||||
historyPath?: string
|
||||
runtimeDir?: string
|
||||
packagedAppVersion?: string | null
|
||||
recoveryOnly?: boolean
|
||||
respawn?: (reason: DaemonRespawnReason) => Promise<void | (() => void)>
|
||||
}
|
||||
|
||||
@@ -71,8 +72,15 @@ export type DaemonIdentityChangeEvent = {
|
||||
current: DaemonEndpointIdentity
|
||||
}
|
||||
|
||||
export type DaemonIdleRetirementResult =
|
||||
| { state: 'retiring' }
|
||||
| { state: 'busy'; liveSessions: number | null; admissionReopened?: true }
|
||||
| { state: 'unsupported' }
|
||||
| { state: 'unverifiable' }
|
||||
|
||||
export abstract class DaemonPtyRuntimeState {
|
||||
readonly protocolVersion: number
|
||||
readonly recoveryOnly: boolean
|
||||
protected socketPath: string
|
||||
protected tokenPath: string
|
||||
protected pidPath: string | null
|
||||
@@ -92,6 +100,9 @@ export abstract class DaemonPtyRuntimeState {
|
||||
protected packagedAppVersion: string | null
|
||||
protected pendingRespawnAdoptionRelease: (() => void) | null = null
|
||||
protected respawnAdoptionClosed = false
|
||||
protected idleRetirementAdmissionClosed = false
|
||||
protected idleRetirementState: 'open' | 'checking' | 'retiring' | 'unverifiable' = 'open'
|
||||
protected idleRetirementPromise: Promise<DaemonIdleRetirementResult> | null = null
|
||||
protected respawnPromise: Promise<void> | null = null
|
||||
protected staleBundleReplacementPromise: Promise<void> | null = null
|
||||
protected writeRecoveryPromise: Promise<void> | null = null
|
||||
@@ -190,6 +201,7 @@ export abstract class DaemonPtyRuntimeState {
|
||||
|
||||
constructor(opts: DaemonPtyAdapterOptions) {
|
||||
this.protocolVersion = opts.protocolVersion ?? PROTOCOL_VERSION
|
||||
this.recoveryOnly = opts.recoveryOnly === true
|
||||
this.socketPath = opts.socketPath
|
||||
this.tokenPath = opts.tokenPath
|
||||
this.pidPath = opts.pidPath ?? null
|
||||
@@ -209,7 +221,7 @@ export abstract class DaemonPtyRuntimeState {
|
||||
})
|
||||
this.historyManager = opts.historyPath ? new HistoryManager(opts.historyPath) : null
|
||||
this.historyReader = opts.historyPath ? new HistoryReader(opts.historyPath) : null
|
||||
this.respawnFn = opts.respawn ?? null
|
||||
this.respawnFn = this.recoveryOnly ? null : (opts.respawn ?? null)
|
||||
this.runtimeDir = opts.runtimeDir ?? opts.profileScope ?? null
|
||||
this.packagedAppVersion = opts.packagedAppVersion ?? null
|
||||
this.supportsCheckpoints = this.protocolVersion >= 4
|
||||
|
||||
@@ -22,10 +22,15 @@ import { resolveSafePtyDefaultCwd } from '../providers/pty-default-cwd'
|
||||
import { resolveUnixShellPath } from '../providers/local-pty-utils'
|
||||
import type { PtySpawnOptions, PtySpawnResult } from '../providers/types'
|
||||
import { injectHistoryEnv, injectWslFishHistoryEnv, logHistoryInjection } from '../terminal-history'
|
||||
import { assertDaemonRecoverySpawnAdmission } from './daemon-recovery-spawn-admission'
|
||||
import { addWslEnvKeys } from '../wsl-env'
|
||||
|
||||
export abstract class DaemonPtySessionSpawn extends DaemonPtySpawnResult {
|
||||
async spawn(opts: PtySpawnOptions): Promise<PtySpawnResult> {
|
||||
assertDaemonRecoverySpawnAdmission(this.recoveryOnly, this.protocolVersion, opts)
|
||||
if (this.idleRetirementAdmissionClosed) {
|
||||
throw new Error('Terminal daemon is decommissioning')
|
||||
}
|
||||
const spawnOpts = this.withHistoryIsolation(opts)
|
||||
const sessionId = spawnOpts.sessionId ?? mintPtySessionId(spawnOpts.worktreeId)
|
||||
const operation: PendingDaemonSpawnOperation = {
|
||||
@@ -105,6 +110,10 @@ export abstract class DaemonPtySessionSpawn extends DaemonPtySpawnResult {
|
||||
operation: PendingDaemonSpawnOperation,
|
||||
historyRecovery: HistoryRecoveryContext
|
||||
): Promise<PtySpawnResult> {
|
||||
assertDaemonRecoverySpawnAdmission(this.recoveryOnly, this.protocolVersion, opts)
|
||||
if (this.idleRetirementAdmissionClosed) {
|
||||
throw new Error('Terminal daemon is decommissioning')
|
||||
}
|
||||
if (
|
||||
opts.agentSessionEnsure &&
|
||||
this.protocolVersion < AGENT_SESSION_CLAIM_DAEMON_PROTOCOL_VERSION
|
||||
|
||||
@@ -0,0 +1,103 @@
|
||||
import { rmSync } from 'node:fs'
|
||||
import { afterEach, beforeEach, expect, it, vi } from 'vitest'
|
||||
import { DaemonPtyAdapter } from './daemon-pty-adapter'
|
||||
import {
|
||||
createMockSubprocess,
|
||||
startDaemonAdapterHarness,
|
||||
waitFor,
|
||||
type DaemonAdapterHarness
|
||||
} from './daemon-pty-adapter-test-harness'
|
||||
import { STABLE_PANE_ATTACH_ONLY_DAEMON_PROTOCOL_VERSION } from './daemon-protocol-version'
|
||||
import { DaemonPtyRouter } from './daemon-pty-router'
|
||||
|
||||
let harness: DaemonAdapterHarness
|
||||
let recovery: DaemonPtyAdapter
|
||||
let subprocess: ReturnType<typeof createMockSubprocess>
|
||||
const spawn = vi.fn(() => subprocess)
|
||||
const respawn = vi.fn(async () => {})
|
||||
|
||||
beforeEach(async () => {
|
||||
spawn.mockClear()
|
||||
respawn.mockClear()
|
||||
subprocess = createMockSubprocess()
|
||||
harness = await startDaemonAdapterHarness(spawn)
|
||||
await harness.adapter.spawn({ sessionId: 'existing', cols: 80, rows: 24 })
|
||||
recovery = new DaemonPtyAdapter({
|
||||
socketPath: harness.socketPath,
|
||||
tokenPath: harness.tokenPath,
|
||||
recoveryOnly: true,
|
||||
respawn
|
||||
})
|
||||
})
|
||||
|
||||
afterEach(async () => {
|
||||
recovery?.dispose()
|
||||
harness.adapter.dispose()
|
||||
await harness.server.shutdown()
|
||||
rmSync(harness.dir, { recursive: true, force: true })
|
||||
})
|
||||
|
||||
it('reattaches and controls existing work without admitting a new process', async () => {
|
||||
await expect(
|
||||
recovery.spawn({ sessionId: 'existing', attachOnly: true, cols: 80, rows: 24 })
|
||||
).resolves.toMatchObject({ id: 'existing', isReattach: true })
|
||||
recovery.write('existing', 'still live\n')
|
||||
await waitFor(() => subprocess.write.mock.calls.length > 0)
|
||||
expect(subprocess.write).toHaveBeenCalledWith('still live\n')
|
||||
await expect(recovery.spawn({ sessionId: 'fresh', cols: 80, rows: 24 })).rejects.toThrow(
|
||||
'managed-stop recovery'
|
||||
)
|
||||
await expect(
|
||||
recovery.spawn({ sessionId: 'missing', attachOnly: true, cols: 80, rows: 24 })
|
||||
).rejects.toThrow()
|
||||
expect(spawn).toHaveBeenCalledOnce()
|
||||
expect(respawn).not.toHaveBeenCalled()
|
||||
})
|
||||
|
||||
it('does not certify fresh admission after a confirmed native stop refusal', async () => {
|
||||
await expect(recovery.requestIdleRetirement()).resolves.toEqual({
|
||||
state: 'busy',
|
||||
liveSessions: 1
|
||||
})
|
||||
await expect(
|
||||
recovery.spawn({ sessionId: 'existing', attachOnly: true, cols: 80, rows: 24 })
|
||||
).resolves.toMatchObject({ isReattach: true })
|
||||
await expect(recovery.spawn({ cols: 80, rows: 24 })).rejects.toThrow('managed-stop recovery')
|
||||
expect(spawn).toHaveBeenCalledOnce()
|
||||
})
|
||||
|
||||
it('does not invoke a supplied replacement launcher after endpoint loss', async () => {
|
||||
await recovery.listProcesses()
|
||||
await harness.server.shutdown()
|
||||
await expect(
|
||||
recovery.spawn({ sessionId: 'existing', attachOnly: true, cols: 80, rows: 24 })
|
||||
).rejects.toThrow()
|
||||
expect(respawn).not.toHaveBeenCalled()
|
||||
})
|
||||
|
||||
it('keeps recovery-only admission through routed inventory refusal', async () => {
|
||||
const router = new DaemonPtyRouter({ current: recovery, legacy: [] })
|
||||
await expect(router.requestIdleRetirement()).resolves.toEqual({ state: 'busy', liveSessions: 1 })
|
||||
await expect(
|
||||
router.spawn({ sessionId: 'existing', attachOnly: true, cols: 80, rows: 24 })
|
||||
).resolves.toMatchObject({ isReattach: true })
|
||||
await expect(router.spawn({ cols: 80, rows: 24 })).rejects.toThrow('managed-stop recovery')
|
||||
expect(spawn).toHaveBeenCalledOnce()
|
||||
})
|
||||
|
||||
it('rejects legacy attach emulation before contacting the daemon', async () => {
|
||||
const legacy = new DaemonPtyAdapter({
|
||||
socketPath: harness.socketPath,
|
||||
tokenPath: harness.tokenPath,
|
||||
protocolVersion: STABLE_PANE_ATTACH_ONLY_DAEMON_PROTOCOL_VERSION - 1,
|
||||
recoveryOnly: true
|
||||
})
|
||||
try {
|
||||
await expect(
|
||||
legacy.spawn({ sessionId: 'existing', attachOnly: true, cols: 80, rows: 24 })
|
||||
).rejects.toThrow('managed-stop recovery')
|
||||
expect(spawn).toHaveBeenCalledOnce()
|
||||
} finally {
|
||||
legacy.dispose()
|
||||
}
|
||||
})
|
||||
@@ -0,0 +1,10 @@
|
||||
import { rebindLocalProviderListeners } from '../ipc/pty'
|
||||
import { getDaemonRuntimeDir, getDaemonHistoryDir } from './daemon-launch-paths'
|
||||
import { createDaemonRecoveryProvider } from './daemon-recovery-provider'
|
||||
import { installDaemonProvider } from './daemon-provider-state'
|
||||
|
||||
export function initDaemonRecoveryProvider(): void {
|
||||
const provider = createDaemonRecoveryProvider(getDaemonRuntimeDir(), getDaemonHistoryDir())
|
||||
installDaemonProvider(null, provider)
|
||||
rebindLocalProviderListeners()
|
||||
}
|
||||
@@ -0,0 +1,54 @@
|
||||
import { copyFileSync, readFileSync, rmSync, writeFileSync } from 'node:fs'
|
||||
import { join } from 'node:path'
|
||||
import { afterEach, beforeEach, expect, it } from 'vitest'
|
||||
import { createDaemonRecoveryProvider } from './daemon-recovery-provider'
|
||||
import { getDaemonPidPath, getDaemonTokenPath } from './daemon-spawner'
|
||||
import { PREVIOUS_DAEMON_PROTOCOL_VERSIONS } from './types'
|
||||
import {
|
||||
createMockSubprocess,
|
||||
startDaemonAdapterHarness,
|
||||
type DaemonAdapterHarness
|
||||
} from './daemon-pty-adapter-test-harness'
|
||||
import type { DaemonPtyRouter } from './daemon-pty-router'
|
||||
|
||||
let harness: DaemonAdapterHarness
|
||||
let provider: DaemonPtyRouter | undefined
|
||||
beforeEach(async () => {
|
||||
harness = await startDaemonAdapterHarness(() => createMockSubprocess())
|
||||
copyFileSync(harness.tokenPath, getDaemonTokenPath(harness.dir))
|
||||
})
|
||||
afterEach(async () => {
|
||||
provider?.dispose()
|
||||
provider = undefined
|
||||
harness.adapter.dispose()
|
||||
await harness.server.shutdown()
|
||||
rmSync(harness.dir, { recursive: true, force: true })
|
||||
})
|
||||
|
||||
it('recovers live terminals without pruning folder or unknown workspace sessions', async () => {
|
||||
await harness.adapter.spawn({ sessionId: 'folder-terminal', cols: 80, rows: 24 })
|
||||
provider = createDaemonRecoveryProvider(harness.dir, join(harness.dir, 'history'))
|
||||
expect(provider.getAllAdapters().every((entry) => entry.recoveryOnly)).toBe(true)
|
||||
await expect(provider.reconcileOnStartup(new Set())).resolves.toEqual({
|
||||
alive: ['folder-terminal'],
|
||||
killed: []
|
||||
})
|
||||
await expect(
|
||||
provider.spawn({ sessionId: 'folder-terminal', attachOnly: true, cols: 80, rows: 24 })
|
||||
).resolves.toMatchObject({ isReattach: true })
|
||||
await expect(provider.spawn({ cols: 80, rows: 24 })).rejects.toThrow('managed-stop recovery')
|
||||
})
|
||||
|
||||
it('retains unreachable legacy generations and their credentials', async () => {
|
||||
const version = PREVIOUS_DAEMON_PROTOCOL_VERSIONS[0]
|
||||
const tokenPath = getDaemonTokenPath(harness.dir, version)
|
||||
const pidPath = getDaemonPidPath(harness.dir, version)
|
||||
writeFileSync(tokenPath, 'retained-secret')
|
||||
writeFileSync(pidPath, '{unreadable pid')
|
||||
provider = createDaemonRecoveryProvider(harness.dir, join(harness.dir, 'history'))
|
||||
const legacy = provider.getAllAdapters().find((entry) => entry.protocolVersion === version)
|
||||
expect(legacy?.recoveryOnly).toBe(true)
|
||||
await expect(legacy!.listSessions()).rejects.toThrow()
|
||||
expect(readFileSync(tokenPath, 'utf8')).toBe('retained-secret')
|
||||
expect(readFileSync(pidPath, 'utf8')).toBe('{unreadable pid')
|
||||
})
|
||||
@@ -0,0 +1,41 @@
|
||||
import { lstatSync } from 'node:fs'
|
||||
import { DaemonPtyAdapter } from './daemon-pty-adapter'
|
||||
import { DaemonPtyRouter } from './daemon-pty-router'
|
||||
import { getDaemonPidPath, getDaemonSocketPath, getDaemonTokenPath } from './daemon-spawner'
|
||||
import { PREVIOUS_DAEMON_PROTOCOL_VERSIONS, PROTOCOL_VERSION } from './types'
|
||||
|
||||
function hasEndpointEvidence(path: string): boolean {
|
||||
try {
|
||||
lstatSync(path)
|
||||
return true
|
||||
} catch (error) {
|
||||
// Unreadable evidence must keep the generation represented as unverifiable.
|
||||
return !(error instanceof Error && 'code' in error && error.code === 'ENOENT')
|
||||
}
|
||||
}
|
||||
|
||||
export function createDaemonRecoveryProvider(
|
||||
runtimeDir: string,
|
||||
historyPath: string
|
||||
): DaemonPtyRouter {
|
||||
const create = (protocolVersion: number): DaemonPtyAdapter =>
|
||||
new DaemonPtyAdapter({
|
||||
socketPath: getDaemonSocketPath(runtimeDir, protocolVersion),
|
||||
tokenPath: getDaemonTokenPath(runtimeDir, protocolVersion),
|
||||
pidPath: getDaemonPidPath(runtimeDir, protocolVersion),
|
||||
profileScope: runtimeDir,
|
||||
runtimeDir,
|
||||
historyPath,
|
||||
protocolVersion,
|
||||
recoveryOnly: true
|
||||
})
|
||||
const legacy = PREVIOUS_DAEMON_PROTOCOL_VERSIONS.filter((version) =>
|
||||
[
|
||||
getDaemonPidPath(runtimeDir, version),
|
||||
getDaemonTokenPath(runtimeDir, version),
|
||||
...(process.platform === 'win32' ? [] : [getDaemonSocketPath(runtimeDir, version)])
|
||||
].some(hasEndpointEvidence)
|
||||
).map(create)
|
||||
// Represent the current endpoint even when absent; missing contact is not an empty census.
|
||||
return new DaemonPtyRouter({ current: create(PROTOCOL_VERSION), legacy })
|
||||
}
|
||||
@@ -0,0 +1,21 @@
|
||||
import type { PtySpawnOptions } from '../providers/types'
|
||||
import { STABLE_PANE_ATTACH_ONLY_DAEMON_PROTOCOL_VERSION } from './daemon-protocol-version'
|
||||
|
||||
export function assertDaemonRecoverySpawnAdmission(
|
||||
recoveryOnly: boolean,
|
||||
protocolVersion: number,
|
||||
opts: PtySpawnOptions
|
||||
): void {
|
||||
if (!recoveryOnly) {
|
||||
return
|
||||
}
|
||||
// Legacy attach emulation can create before rejecting the result.
|
||||
if (
|
||||
opts.attachOnly !== true ||
|
||||
!opts.sessionId ||
|
||||
opts.agentSessionEnsure ||
|
||||
protocolVersion < STABLE_PANE_ATTACH_ONLY_DAEMON_PROTOCOL_VERSION
|
||||
) {
|
||||
throw new Error('Terminal daemon admission is fenced for managed-stop recovery')
|
||||
}
|
||||
}
|
||||
@@ -10,6 +10,7 @@ import type { DaemonSessionBackgroundRouting } from './daemon-session-background
|
||||
import { recordDaemonStreamBacklogEvent } from './daemon-stream-backlog-probe'
|
||||
import type { DaemonStreamDataBatcher } from './daemon-stream-data-batcher'
|
||||
import type { DaemonTerminalAdmission } from './daemon-terminal-admission'
|
||||
import { readDaemonHealthIdentity } from './daemon-health-identity'
|
||||
import type { TerminalHistorySeedTransferRegistry } from './terminal-history-seed-transfer-registry'
|
||||
import type { TerminalHost } from './terminal-host'
|
||||
import { SessionNotFoundError, type DaemonRequest } from './types'
|
||||
@@ -156,7 +157,7 @@ export class DaemonRequestRouter {
|
||||
return { health: await readCurrentProcessMacSystemResolverHealth() }
|
||||
case 'ptySpawnHealth':
|
||||
await this.options.ptySpawnHealthCheck()
|
||||
return { healthy: true }
|
||||
return { healthy: true, ...readDaemonHealthIdentity() }
|
||||
case 'shutdown':
|
||||
return this.shutdown(clientId, request.id, request.payload.killSessions)
|
||||
}
|
||||
|
||||
@@ -0,0 +1,49 @@
|
||||
import { expect, it, vi } from 'vitest'
|
||||
import { DaemonRouterRetirement } from './daemon-router-retirement'
|
||||
import { createAdapter } from './daemon-pty-router-test-fixture'
|
||||
import { PROTOCOL_VERSION } from './types'
|
||||
|
||||
it.each(['inventory', 'protocol', 'spawn', 'live'] as const)(
|
||||
'does not reopen admission on a %s retry after partial retirement',
|
||||
async (failure) => {
|
||||
const current = createAdapter('current', [], undefined, PROTOCOL_VERSION)
|
||||
const legacy = createAdapter('legacy', [], undefined, PROTOCOL_VERSION)
|
||||
let adapters = [current, legacy]
|
||||
const retirement = new DaemonRouterRetirement(() => adapters)
|
||||
vi.mocked(legacy.requestIdleRetirement).mockResolvedValueOnce({
|
||||
state: 'busy',
|
||||
liveSessions: 0
|
||||
})
|
||||
await expect(retirement.requestIdleRetirement()).resolves.toEqual({ state: 'unverifiable' })
|
||||
expect(retirement.admissionClosed).toBe(true)
|
||||
if (failure === 'inventory') {
|
||||
vi.mocked(current.listSessions).mockRejectedValueOnce(new Error('lost connection'))
|
||||
} else if (failure === 'protocol') {
|
||||
adapters = [createAdapter('old', [], undefined, 23)]
|
||||
} else if (failure === 'spawn') {
|
||||
retirement.spawnInFlight = 1
|
||||
} else {
|
||||
adapters = [createAdapter('live', ['existing'], undefined, PROTOCOL_VERSION)]
|
||||
}
|
||||
expect(await retirement.requestIdleRetirement()).not.toHaveProperty('admissionReopened')
|
||||
expect(retirement.admissionClosed).toBe(true)
|
||||
}
|
||||
)
|
||||
|
||||
it('does not reopen when every native result is busy without reopening proof', async () => {
|
||||
const current = createAdapter('current', [], undefined, PROTOCOL_VERSION)
|
||||
vi.mocked(current.requestIdleRetirement).mockResolvedValue({ state: 'busy', liveSessions: 0 })
|
||||
const retirement = new DaemonRouterRetirement(() => [current])
|
||||
await expect(retirement.requestIdleRetirement()).resolves.toEqual({ state: 'unverifiable' })
|
||||
expect(retirement.admissionClosed).toBe(true)
|
||||
})
|
||||
|
||||
it('keeps the fence through a lost native reply and failed retry inventory', async () => {
|
||||
const current = createAdapter('current', [], undefined, PROTOCOL_VERSION)
|
||||
vi.mocked(current.requestIdleRetirement).mockRejectedValueOnce(new Error('lost stop reply'))
|
||||
const retirement = new DaemonRouterRetirement(() => [current])
|
||||
await expect(retirement.requestIdleRetirement()).rejects.toThrow('lost stop reply')
|
||||
vi.mocked(current.listSessions).mockRejectedValueOnce(new Error('lost connection'))
|
||||
await expect(retirement.requestIdleRetirement()).resolves.toEqual({ state: 'unverifiable' })
|
||||
expect(retirement.admissionClosed).toBe(true)
|
||||
})
|
||||
@@ -0,0 +1,74 @@
|
||||
import type { DaemonPtyAdapter } from './daemon-pty-adapter'
|
||||
import type { DaemonIdleRetirementResult } from './daemon-pty-runtime-state'
|
||||
import { CLEAN_DISCONNECT_PROTOCOL_VERSION } from './types'
|
||||
|
||||
export class DaemonRouterRetirement {
|
||||
admissionClosed = false
|
||||
spawnInFlight = 0
|
||||
private retirementAttempted = false
|
||||
private idleRetirementPromise: Promise<DaemonIdleRetirementResult> | null = null
|
||||
|
||||
constructor(private readonly allAdapters: () => DaemonPtyAdapter[]) {}
|
||||
|
||||
async requestIdleRetirement(): Promise<DaemonIdleRetirementResult> {
|
||||
if (this.idleRetirementPromise) {
|
||||
return this.idleRetirementPromise
|
||||
}
|
||||
this.admissionClosed = true
|
||||
const request = this.finishIdleRetirementRequest().finally(() => {
|
||||
if (this.idleRetirementPromise === request) {
|
||||
this.idleRetirementPromise = null
|
||||
}
|
||||
})
|
||||
this.idleRetirementPromise = request
|
||||
return request
|
||||
}
|
||||
|
||||
private async finishIdleRetirementRequest(): Promise<DaemonIdleRetirementResult> {
|
||||
const adapters = this.allAdapters()
|
||||
if (this.spawnInFlight > 0) {
|
||||
this.admissionClosed = this.retirementAttempted
|
||||
return { state: 'busy', liveSessions: null }
|
||||
}
|
||||
if (adapters.some((adapter) => adapter.protocolVersion < CLEAN_DISCONNECT_PROTOCOL_VERSION)) {
|
||||
this.admissionClosed = this.retirementAttempted
|
||||
return { state: 'unsupported' }
|
||||
}
|
||||
const inventories = await Promise.allSettled(adapters.map((adapter) => adapter.listSessions()))
|
||||
if (inventories.some((inventory) => inventory.status === 'rejected')) {
|
||||
this.admissionClosed = this.retirementAttempted
|
||||
return { state: 'unverifiable' }
|
||||
}
|
||||
const liveSessions = inventories.reduce(
|
||||
(count, inventory) =>
|
||||
count +
|
||||
(inventory.status === 'fulfilled'
|
||||
? inventory.value.filter((session) => session.isAlive).length
|
||||
: 0),
|
||||
0
|
||||
)
|
||||
if (liveSessions > 0) {
|
||||
this.admissionClosed = this.retirementAttempted
|
||||
return {
|
||||
state: 'busy',
|
||||
liveSessions,
|
||||
...(!this.retirementAttempted && !adapters.some((adapter) => adapter.recoveryOnly)
|
||||
? { admissionReopened: true as const }
|
||||
: {})
|
||||
}
|
||||
}
|
||||
this.retirementAttempted = true
|
||||
const results = await Promise.all(adapters.map((adapter) => adapter.requestIdleRetirement()))
|
||||
if (results.every((result) => result.state === 'retiring')) {
|
||||
return { state: 'retiring' }
|
||||
}
|
||||
const refusedLiveSessions = results.reduce(
|
||||
(count, result) => count + (result.state === 'busy' ? (result.liveSessions ?? 0) : 0),
|
||||
0
|
||||
)
|
||||
if (refusedLiveSessions > 0) {
|
||||
return { state: 'busy', liveSessions: refusedLiveSessions }
|
||||
}
|
||||
return { state: 'unverifiable' }
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,37 @@
|
||||
import type { DaemonPtyAdapter } from './daemon-pty-adapter'
|
||||
import type { DaemonSessionOwnerResolver } from './daemon-session-owner-resolution'
|
||||
|
||||
export async function reconcileDaemonRouterSessions(
|
||||
adapters: readonly DaemonPtyAdapter[],
|
||||
ownerResolver: DaemonSessionOwnerResolver<DaemonPtyAdapter>,
|
||||
validWorktreeIds: Set<string>
|
||||
): Promise<{ alive: string[]; killed: string[] }> {
|
||||
const alive: string[] = []
|
||||
const killed: string[] = []
|
||||
const aliveProviders = new Map<string, Set<DaemonPtyAdapter>>()
|
||||
for (const adapter of adapters) {
|
||||
const result = await adapter.reconcileOnStartup(validWorktreeIds)
|
||||
// Why: daemon startup can reconcile many restored sessions; spreading
|
||||
// those arrays into push can exceed JavaScript's argument limit.
|
||||
for (const id of result.alive) {
|
||||
alive.push(id)
|
||||
}
|
||||
for (const id of result.killed) {
|
||||
killed.push(id)
|
||||
}
|
||||
for (const id of result.alive) {
|
||||
const providers = aliveProviders.get(id) ?? new Set<DaemonPtyAdapter>()
|
||||
providers.add(adapter)
|
||||
aliveProviders.set(id, providers)
|
||||
}
|
||||
}
|
||||
for (const id of new Set([...alive, ...killed])) {
|
||||
const providers = aliveProviders.get(id)
|
||||
if (providers?.size === 1) {
|
||||
ownerResolver.recordRoute(id, providers.values().next().value!)
|
||||
} else {
|
||||
ownerResolver.forgetRoute(id)
|
||||
}
|
||||
}
|
||||
return { alive, killed }
|
||||
}
|
||||
Reference in New Issue
Block a user