diff --git a/src/main/runtime/orchestration/db/worker-terminal/worker-terminal-release.ts b/src/main/runtime/orchestration/db/worker-terminal/worker-terminal-release.ts index 832b32e5c45..5e96249cdde 100644 --- a/src/main/runtime/orchestration/db/worker-terminal/worker-terminal-release.ts +++ b/src/main/runtime/orchestration/db/worker-terminal/worker-terminal-release.ts @@ -1,4 +1,8 @@ -import { WORKER_SETTLED_STATES } from '../../worker-terminal-ownership' +import { + decideWorkerTerminalRelease, + WORKER_SETTLED_STATES, + WORKER_TERMINAL_RELEASABLE_ROW_SQL +} from '../../worker-terminal-ownership' import type { WorkerTerminalResourceRow, WorkerTerminalRetainedReason @@ -51,7 +55,8 @@ export function requestWorkerTerminalRelease( ? { disposition: 'retained', resource: transferred, reason: 'ownership_transferred' } : { disposition: 'retained', resource: null, reason: 'no_owned_resource' } } - if (resource.release_state === 'released' || resource.ownership_state === 'released') { + const decision = decideWorkerTerminalRelease(resource) + if (decision.action === 'already_released') { this.db.exec('COMMIT') return { disposition: 'already_released', resource } } @@ -59,21 +64,9 @@ export function requestWorkerTerminalRelease( this.db.exec('COMMIT') return { disposition: 'retained', resource, reason: 'identity_unproven' } } - if (resource.ownership_state === 'external') { + if (decision.action === 'retained') { this.db.exec('COMMIT') - return { - disposition: 'retained', - resource, - reason: (resource.retained_reason as WorkerTerminalRetainedReason) ?? 'external_terminal' - } - } - if (resource.ownership_state === 'user_owned') { - this.db.exec('COMMIT') - return { disposition: 'retained', resource, reason: 'user_takeover' } - } - if (resource.ownership_state === 'transferred') { - this.db.exec('COMMIT') - return { disposition: 'retained', resource, reason: 'ownership_transferred' } + return { disposition: 'retained', resource, reason: decision.reason } } if (resource.release_state === 'retained' && resource.retained_reason === 'user_requested') { this.db.prepare('DELETE FROM worker_terminal_archives WHERE dispatch_id = ?').run(dispatchId) @@ -88,7 +81,7 @@ export function requestWorkerTerminalRelease( retained_reason = NULL, release_requested_at = COALESCE(release_requested_at, datetime('now')), release_error = NULL, updated_at = datetime('now') - WHERE id = ? AND release_state IN ('not_requested', 'retained', 'requested', 'releasing', 'unknown')` + WHERE id = ? AND ${WORKER_TERMINAL_RELEASABLE_ROW_SQL}` ) .run(resource.id) this.db.exec('COMMIT') @@ -108,7 +101,6 @@ export function settleDeadWorkerTerminalRelease( requestingDispatchId: string resourceId: string processIncarnation: string - requireArchive?: boolean } ): | { disposition: 'released'; resource: WorkerTerminalResourceRow } @@ -132,20 +124,16 @@ export function settleDeadWorkerTerminalRelease( const ownerSettled = Boolean(owner && WORKER_SETTLED_STATES.includes(owner.state)) // A positive process-exit verdict only proves the exact process is gone; release is // terminal cleanup and must also have a durable output archive to preserve worker evidence. - const archive = params.requireArchive - ? this.getWorkerTerminalArchive(resource.owner_dispatch_id) - : undefined + const archive = this.getWorkerTerminalArchive(resource.owner_dispatch_id) if ( !priorOwners || !requesterRelated || !requesterSettled || !ownerSettled || resource.process_incarnation !== params.processIncarnation || - (params.requireArchive && (!archive || archive.resource_id !== resource.id)) || - resource.ownership_state === 'released' || - !['not_requested', 'retained', 'requested', 'releasing', 'unknown'].includes( - resource.release_state - ) + !archive || + archive.resource_id !== resource.id || + decideWorkerTerminalRelease(resource).action !== 'proceed' ) { this.db.exec('COMMIT') return { disposition: 'retained', resource } @@ -157,8 +145,7 @@ export function settleDeadWorkerTerminalRelease( release_requested_at = COALESCE(release_requested_at, datetime('now')), release_completed_at = datetime('now'), release_error = NULL, updated_at = datetime('now') - WHERE id = ? AND process_incarnation = ? AND ownership_state != 'released' - AND release_state IN ('not_requested', 'retained', 'requested', 'releasing', 'unknown')` + WHERE id = ? AND process_incarnation = ? AND ${WORKER_TERMINAL_RELEASABLE_ROW_SQL}` ) .run(params.resourceId, params.processIncarnation) const released = this.getWorkerTerminalResource(params.resourceId) as WorkerTerminalResourceRow diff --git a/src/main/runtime/orchestration/worker-terminal-ownership.ts b/src/main/runtime/orchestration/worker-terminal-ownership.ts index 43ec7418e94..a4a3d729095 100644 --- a/src/main/runtime/orchestration/worker-terminal-ownership.ts +++ b/src/main/runtime/orchestration/worker-terminal-ownership.ts @@ -113,3 +113,35 @@ export function deriveWorkerTerminalListState(params: { ? 'retained' : 'active' } + +export type WorkerTerminalReleaseDecision = + | { action: 'already_released' } + | { action: 'retained'; reason: WorkerTerminalRetainedReason } + | { action: 'proceed' } + +// The single (ownership_state, release_state) -> action table. Both release guards read it, so a +// resource the dispatch no longer owns can never be settled as released down either path. +export function decideWorkerTerminalRelease( + resource: Pick +): WorkerTerminalReleaseDecision { + if (resource.release_state === 'released' || resource.ownership_state === 'released') { + return { action: 'already_released' } + } + switch (resource.ownership_state) { + case 'external': + return { + action: 'retained', + reason: (resource.retained_reason as WorkerTerminalRetainedReason) ?? 'external_terminal' + } + case 'user_owned': + return { action: 'retained', reason: 'user_takeover' } + case 'transferred': + return { action: 'retained', reason: 'ownership_transferred' } + case 'owned': + return { action: 'proceed' } + } +} + +/** SQL form of the table's `proceed` arm, for the compare-and-set race guard on the same row. */ +export const WORKER_TERMINAL_RELEASABLE_ROW_SQL = + "ownership_state = 'owned' AND release_state <> 'released'" diff --git a/src/main/runtime/rpc/methods/orchestration-worker-release-completion.ts b/src/main/runtime/rpc/methods/orchestration-worker-release-completion.ts index 584d0c5b4ee..bc1315e7008 100644 --- a/src/main/runtime/rpc/methods/orchestration-worker-release-completion.ts +++ b/src/main/runtime/rpc/methods/orchestration-worker-release-completion.ts @@ -142,8 +142,7 @@ async function completeWorkerTerminalReleaseOnce( const reconciled = db.settleDeadWorkerTerminalRelease({ requestingDispatchId: dispatchId, resourceId: resource.id, - processIncarnation: resource.process_incarnation, - requireArchive: true + processIncarnation: resource.process_incarnation }) if (reconciled.disposition === 'released') { runtime.notifyMessageArrived(`dispatch:${dispatchId}`, 'status') diff --git a/src/main/runtime/rpc/methods/orchestration-worker-release-inventory.test.ts b/src/main/runtime/rpc/methods/orchestration-worker-release-inventory.test.ts index 3f2073c13fa..913ff5223e3 100644 --- a/src/main/runtime/rpc/methods/orchestration-worker-release-inventory.test.ts +++ b/src/main/runtime/rpc/methods/orchestration-worker-release-inventory.test.ts @@ -34,7 +34,7 @@ describe('orchestration worker release inventory', () => { expect(h.runtime.closeTerminal).toHaveBeenCalledWith('term_reminted') }) - it('reconciles dead transferred ownership after the current owner settles', async () => { + it('refuses to settle dead transferred ownership with no durable archive', async () => { h.setup() const first = await h.startSettledWorker('succeeded') const second = await h.startWorker({ terminal: 'term_reminted' }) @@ -43,16 +43,19 @@ describe('orchestration worker release inventory', () => { await expect( h.call('orchestration.workerRelease', { dispatch: first.dispatchId }) - ).resolves.toMatchObject({ state: 'released', processAction: 'none' }) + ).resolves.toMatchObject({ + state: 'retained', + reason: 'ownership_transferred', + processAction: 'none' + }) expect(h.inspectProcessLiveness).toHaveBeenCalledWith( 'runtime_test:term_worker:1', JSON.stringify({ kind: 'local', hostId: 'local' }) ) expect(h.runtime.closeTerminal).not.toHaveBeenCalled() - expect(h.db.getWorkerTerminalResourceByOwner(second.dispatchId)).toMatchObject({ - ownership_state: 'released', - release_state: 'released' - }) + expect(h.db.getWorkerTerminalResourceByOwner(second.dispatchId)?.release_state).not.toBe( + 'released' + ) }) it('rejects exact reuse after release intent instead of closing the new worker', async () => { diff --git a/src/main/runtime/rpc/methods/orchestration-worker-release-ownership-guard.test.ts b/src/main/runtime/rpc/methods/orchestration-worker-release-ownership-guard.test.ts new file mode 100644 index 00000000000..8f6e27849b8 --- /dev/null +++ b/src/main/runtime/rpc/methods/orchestration-worker-release-ownership-guard.test.ts @@ -0,0 +1,69 @@ +import { afterEach, beforeEach, describe, expect, it } from 'vitest' +import { createOrchestrationWorkerReleaseHarness } from './orchestration-worker-release-test-harness' + +describe('workerRelease on a retained resource whose process exited', () => { + const harness = createOrchestrationWorkerReleaseHarness() + beforeEach(() => harness.setup()) + afterEach(() => harness.cleanup()) + + it('does not release a terminal the user took over', async () => { + const { dispatchId } = await harness.startSettledWorker('succeeded') + const takeover = (await harness.call('orchestration.workerTerminalUserInput', { + paneKey: harness.workerPaneKey + })) as { changed: number } + expect(takeover.changed).toBe(1) + expect(harness.db.getWorkerTerminalResourceByOwner(dispatchId)?.ownership_state).toBe( + 'user_owned' + ) + + // The agent process later exits on its own; the user's pane and scrollback remain. + harness.inspectProcessLiveness.mockResolvedValue('exited') + const receipt = (await harness.call('orchestration.workerRelease', { + dispatch: dispatchId + })) as { state: string; reason?: string; archive: unknown } + + expect(receipt.state).toBe('retained') + expect(receipt.reason).toBe('user_takeover') + const after = harness.db.getWorkerTerminalResourceByOwner(dispatchId) + expect(after?.ownership_state).toBe('user_owned') + expect(after?.release_state).not.toBe('released') + }) + + it.each(['transferred', 'external'] as const)( + 'does not release a %s resource on an exited process', + async (ownershipState) => { + const { dispatchId } = await harness.startSettledWorker('succeeded') + const resource = harness.db.getWorkerTerminalResourceByOwner(dispatchId)! + harness.db.db + .prepare('UPDATE worker_terminal_resources SET ownership_state = ? WHERE id = ?') + .run(ownershipState, resource.id) + + harness.inspectProcessLiveness.mockResolvedValue('exited') + const receipt = (await harness.call('orchestration.workerRelease', { + dispatch: dispatchId + })) as { state: string } + + expect(receipt.state).toBe('retained') + const after = harness.db.getWorkerTerminalResourceByOwner(dispatchId) + expect(after?.ownership_state).toBe(ownershipState) + expect(after?.release_state).not.toBe('released') + } + ) + + it('does not mark released without a durable output archive', async () => { + const { dispatchId } = await harness.startWorker() + // Abandon so release reports `identity_unproven` and keeps the still-owned pane. + expect(harness.db.abandonWorkerDispatch(dispatchId).disposition).toBe('abandoned') + expect(harness.db.getWorkerTerminalArchive(dispatchId)).toBeFalsy() + + harness.inspectProcessLiveness.mockResolvedValue('exited') + const receipt = (await harness.call('orchestration.workerRelease', { + dispatch: dispatchId + })) as { state: string } + + expect(receipt.state).toBe('retained') + expect(harness.db.getWorkerTerminalResourceByOwner(dispatchId)?.release_state).not.toBe( + 'released' + ) + }) +}) diff --git a/src/main/runtime/rpc/methods/orchestration-worker-release.test.ts b/src/main/runtime/rpc/methods/orchestration-worker-release.test.ts index 09f73faa70d..34e0d8579f7 100644 --- a/src/main/runtime/rpc/methods/orchestration-worker-release.test.ts +++ b/src/main/runtime/rpc/methods/orchestration-worker-release.test.ts @@ -104,22 +104,26 @@ describe('orchestration worker release', () => { expect(h.runtime.closeTerminal).not.toHaveBeenCalled() }) - it('reconciles a dead external terminal without closing a process', async () => { + it('retains a dead external terminal the orchestration never owned', async () => { h.setup() const { dispatchId } = await h.startSettledWorker('succeeded', { terminal: 'term_worker' }) h.inspectProcessLiveness.mockResolvedValue('exited') await expect( h.call('orchestration.workerRelease', { dispatch: dispatchId }) - ).resolves.toMatchObject({ state: 'released', processAction: 'none' }) + ).resolves.toMatchObject({ + state: 'retained', + reason: 'external_terminal', + processAction: 'none' + }) expect(h.inspectProcessLiveness).toHaveBeenCalledWith( 'runtime_test:term_worker:1', JSON.stringify({ kind: 'local', hostId: 'local' }) ) expect(h.runtime.closeTerminal).not.toHaveBeenCalled() expect(h.db.getWorkerTerminalResourceByOwner(dispatchId)).toMatchObject({ - ownership_state: 'released', - release_state: 'released' + ownership_state: 'external', + release_state: 'not_requested' }) }) @@ -158,7 +162,7 @@ describe('orchestration worker release', () => { expect(h.db.getWorkerTerminalResourceByOwner(dispatchId)?.ownership_state).toBe('user_owned') }) - it('reconciles a dead user-taken-over terminal without closing a process', async () => { + it('keeps a dead user-taken-over terminal in the user takeover', async () => { h.setup() const { dispatchId } = await h.startSettledWorker() await h.call('orchestration.workerTerminalUserInput', { paneKey: h.workerPaneKey }) @@ -166,20 +170,20 @@ describe('orchestration worker release', () => { await expect( h.call('orchestration.workerRelease', { dispatch: dispatchId }) - ).resolves.toMatchObject({ state: 'released', processAction: 'none' }) + ).resolves.toMatchObject({ state: 'retained', reason: 'user_takeover', processAction: 'none' }) expect(h.inspectProcessLiveness).toHaveBeenCalledWith( 'runtime_test:term_worker:1', JSON.stringify({ kind: 'local', hostId: 'local' }) ) expect(h.runtime.closeTerminal).not.toHaveBeenCalled() expect(h.db.getWorkerTerminalResourceByOwner(dispatchId)).toMatchObject({ - ownership_state: 'released', - release_state: 'released' + ownership_state: 'user_owned', + release_state: 'retained' }) }) it.each(['stopped', 'abandoned'] as const)( - 'reconciles a dead %s worker without closing a process', + 'refuses to settle a dead %s worker whose output was never archived', async (state) => { h.setup() const { dispatchId } = await h.startWorker() @@ -193,12 +197,13 @@ describe('orchestration worker release', () => { await expect( h.call('orchestration.workerRelease', { dispatch: dispatchId }) - ).resolves.toMatchObject({ state: 'released', processAction: 'none' }) - expect(h.runtime.closeTerminal).not.toHaveBeenCalled() - expect(h.db.getWorkerTerminalResourceByOwner(dispatchId)).toMatchObject({ - ownership_state: 'released', - release_state: 'released' + ).resolves.toMatchObject({ + state: 'retained', + reason: 'identity_unproven', + processAction: 'none' }) + expect(h.runtime.closeTerminal).not.toHaveBeenCalled() + expect(h.db.getWorkerTerminalResourceByOwner(dispatchId)?.release_state).not.toBe('released') } )