mirror of
https://github.com/stablyai/orca.git
synced 2026-09-30 00:03:15 +00:00
fix(orchestration): never settle a worker terminal the dispatch no longer owns
worker-release's dead-process shortcut called settleDeadWorkerTerminalRelease without requireArchive, and the guard only excluded ownership_state='released', so a user takeover, an external terminal, or a transferred resource was marked released and its output archive was lost. Both release guards now read one (ownership_state, release_state) -> action table, and settlement always requires the durable archive.
This commit is contained in:
@@ -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
|
||||
|
||||
@@ -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<WorkerTerminalResourceRow, 'ownership_state' | 'release_state' | 'retained_reason'>
|
||||
): 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'"
|
||||
|
||||
@@ -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')
|
||||
|
||||
@@ -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 () => {
|
||||
|
||||
@@ -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'
|
||||
)
|
||||
})
|
||||
})
|
||||
@@ -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')
|
||||
}
|
||||
)
|
||||
|
||||
|
||||
Reference in New Issue
Block a user