mirror of
https://github.com/stablyai/orca.git
synced 2026-09-29 08:03:20 +00:00
fix(orchestration): recover released and exited terminal leases
This commit is contained in:
@@ -14,6 +14,7 @@ import { inspectWorkerTerminal } from './worker-observation'
|
||||
import { orchestrationTimestampToMs } from './worker-output'
|
||||
import { archiveSummary } from './worker-terminal-resource-presentation'
|
||||
import { classifyWorkerTerminalCloseError } from './worker-release-close-error'
|
||||
import { workerTerminalLeaseIsCurrent } from './worker-terminal-release-lease'
|
||||
|
||||
export {
|
||||
archiveSummary,
|
||||
@@ -109,16 +110,6 @@ async function completeWorkerTerminalReleaseOnce(
|
||||
archive: archiveSummary(retained)
|
||||
}
|
||||
}
|
||||
if (!workerTerminalLeaseIsCurrent(runtime, db, dispatchId, resource)) {
|
||||
const retained = db.revertWorkerTerminalReleaseToRetained(resource.id, 'identity_unproven')
|
||||
return {
|
||||
dispatchId,
|
||||
state: 'retained',
|
||||
reason: 'identity_unproven',
|
||||
processAction: 'none',
|
||||
archive: archiveSummary(retained)
|
||||
}
|
||||
}
|
||||
if (observation.status === 'missing' || observation.status === 'unattached') {
|
||||
if (args.mode === 'recovery') {
|
||||
// A close can succeed before the process crashes, leaving `releasing` durable state while
|
||||
@@ -172,6 +163,16 @@ async function completeWorkerTerminalReleaseOnce(
|
||||
}
|
||||
}
|
||||
|
||||
if (!workerTerminalLeaseIsCurrent(runtime, db, dispatchId, resource)) {
|
||||
const retained = db.revertWorkerTerminalReleaseToRetained(resource.id, 'identity_unproven')
|
||||
return {
|
||||
dispatchId,
|
||||
state: 'retained',
|
||||
reason: 'identity_unproven',
|
||||
processAction: 'none',
|
||||
archive: archiveSummary(retained)
|
||||
}
|
||||
}
|
||||
const archive = db.getWorkerTerminalArchive(dispatchId)
|
||||
let archiveSource = resource.archive_source as 'transcript' | 'terminal' | null
|
||||
let archiveStatus: WorkerTerminalArchiveStatus | null = resource.archive_status
|
||||
@@ -276,27 +277,6 @@ export function releaseUnknownRecovery(dispatchId: string): string {
|
||||
return `Inspect with: orca orchestration worker-show --dispatch ${dispatchId} --json — then retry worker-release with a fresh request ID (omit --retry-request to let the CLI generate one). Reusing the prior request ID only replays this release_unknown receipt. Never substitute a broad terminal close.`
|
||||
}
|
||||
|
||||
function workerTerminalLeaseIsCurrent(
|
||||
runtime: OrcaRuntimeService,
|
||||
db: OrchestrationDb,
|
||||
dispatchId: string,
|
||||
resource: WorkerTerminalResourceRow
|
||||
): boolean {
|
||||
const worker = db.getWorkerDispatch(dispatchId)
|
||||
const authority = runtime.getOrchestrationDispatchAuthority(resource.terminal_handle)
|
||||
return Boolean(
|
||||
worker?.agent_terminal_handle === resource.terminal_handle &&
|
||||
authority &&
|
||||
resource.host_scope === JSON.stringify(authority.hostScope) &&
|
||||
db.isDispatchProcessCurrent({
|
||||
dispatchId,
|
||||
paneKey: runtime.getTerminalPaneKey(resource.terminal_handle),
|
||||
processIncarnation: runtime.getTerminalProcessIncarnation(resource.terminal_handle)
|
||||
}) &&
|
||||
!db.workerTerminalResourceHasIdentityConflict(resource.id)
|
||||
)
|
||||
}
|
||||
|
||||
function retainedReason(resource: WorkerTerminalResourceRow): WorkerTerminalRetainedReason {
|
||||
if (resource.retained_reason) {
|
||||
return resource.retained_reason as WorkerTerminalRetainedReason
|
||||
|
||||
@@ -168,6 +168,9 @@ describe('orchestration worker release recovery', () => {
|
||||
expect(runtime.closeTerminal).toHaveBeenCalledTimes(1)
|
||||
|
||||
vi.mocked(runtime.showTerminal).mockRejectedValue(new Error('terminal_handle_stale'))
|
||||
vi.mocked(runtime.getOrchestrationDispatchAuthority).mockReturnValue(null)
|
||||
vi.mocked(runtime.getTerminalPaneKey).mockReturnValue(null)
|
||||
vi.mocked(runtime.getTerminalProcessIncarnation).mockReturnValue(null)
|
||||
vi.spyOn(runtime, 'inspectTerminalProcessIncarnationLiveness').mockResolvedValue('exited')
|
||||
|
||||
await expect(reconcileRequestedWorkerTerminalReleases(runtime)).resolves.toMatchObject({
|
||||
@@ -197,6 +200,9 @@ describe('orchestration worker release recovery', () => {
|
||||
expect(db.getWorkerTerminalArchive(dispatchId)).toBeUndefined()
|
||||
|
||||
vi.mocked(runtime.showTerminal).mockRejectedValue(new Error('terminal_handle_stale'))
|
||||
vi.mocked(runtime.getOrchestrationDispatchAuthority).mockReturnValue(null)
|
||||
vi.mocked(runtime.getTerminalPaneKey).mockReturnValue(null)
|
||||
vi.mocked(runtime.getTerminalProcessIncarnation).mockReturnValue(null)
|
||||
vi.spyOn(runtime, 'inspectTerminalProcessIncarnationLiveness').mockResolvedValue('exited')
|
||||
|
||||
await expect(reconcileRequestedWorkerTerminalReleases(runtime)).resolves.toMatchObject({
|
||||
@@ -216,6 +222,8 @@ describe('orchestration worker release recovery', () => {
|
||||
setup()
|
||||
const { dispatchId } = await startSettledWorker()
|
||||
vi.spyOn(runtime, 'getTerminalLivenessVerdict').mockReturnValue({ status: 'exited' })
|
||||
vi.mocked(runtime.getOrchestrationDispatchAuthority).mockRestore()
|
||||
expect(runtime.getOrchestrationDispatchAuthority('term_worker')).toBeNull()
|
||||
vi.mocked(runtime.closeTerminal).mockRejectedValueOnce(new Error('Multiplexer disposed'))
|
||||
|
||||
await expect(
|
||||
@@ -223,6 +231,21 @@ describe('orchestration worker release recovery', () => {
|
||||
).resolves.toMatchObject({ state: 'released', processAction: 'closed_exited_terminal' })
|
||||
})
|
||||
|
||||
it('does not substitute absent launch authority for a positive host exit verdict', async () => {
|
||||
setup()
|
||||
const { dispatchId } = await startSettledWorker()
|
||||
vi.mocked(runtime.getOrchestrationDispatchAuthority).mockRestore()
|
||||
vi.spyOn(runtime, 'getTerminalLivenessVerdict').mockReturnValue({
|
||||
status: 'unverifiable',
|
||||
reason: 'missing_liveness_verdict'
|
||||
})
|
||||
|
||||
await expect(
|
||||
call('orchestration.workerRelease', { dispatch: dispatchId })
|
||||
).resolves.toMatchObject({ state: 'retained', reason: 'identity_unproven' })
|
||||
expect(runtime.closeTerminal).not.toHaveBeenCalled()
|
||||
})
|
||||
|
||||
it('defers instead of settling unknown while inventory is incomplete', async () => {
|
||||
setup()
|
||||
const { dispatchId } = await startSettledWorker()
|
||||
|
||||
@@ -0,0 +1,26 @@
|
||||
import type { OrcaRuntimeService } from '../../../../orca-runtime'
|
||||
import type { OrchestrationDb } from '../../../../orchestration/db'
|
||||
import type { WorkerTerminalResourceRow } from '../../../../orchestration/worker-terminal-ownership'
|
||||
|
||||
export function workerTerminalLeaseIsCurrent(
|
||||
runtime: OrcaRuntimeService,
|
||||
db: OrchestrationDb,
|
||||
dispatchId: string,
|
||||
resource: WorkerTerminalResourceRow
|
||||
): boolean {
|
||||
const worker = db.getWorkerDispatch(dispatchId)
|
||||
const authority = runtime.getOrchestrationDispatchAuthority(resource.terminal_handle)
|
||||
// Exited PTYs retain identity and host evidence but no longer mint launch authority.
|
||||
return Boolean(
|
||||
worker?.agent_terminal_handle === resource.terminal_handle &&
|
||||
(authority
|
||||
? resource.host_scope === JSON.stringify(authority.hostScope)
|
||||
: runtime.getTerminalLivenessVerdict(resource.terminal_handle)?.status === 'exited') &&
|
||||
db.isDispatchProcessCurrent({
|
||||
dispatchId,
|
||||
paneKey: runtime.getTerminalPaneKey(resource.terminal_handle),
|
||||
processIncarnation: runtime.getTerminalProcessIncarnation(resource.terminal_handle)
|
||||
}) &&
|
||||
!db.workerTerminalResourceHasIdentityConflict(resource.id)
|
||||
)
|
||||
}
|
||||
Reference in New Issue
Block a user