mirror of
https://github.com/aipoch/open-science.git
synced 2026-10-02 02:14:40 +08:00
fix(acp): harden Prisma config and continuation admission (#3028)
* fix(deps): harden Prisma config dependency resolution Scope deepmerge-ts 8.0.2 to @prisma/config 6.19.3 to fix recursive graph stack exhaustion without changing Prisma or database engines. Exercise recursive merges and real CJS/MTS config loading through the consumer dependency path. * test(regressions): stabilize scheduled coverage Give the debounced automatic collection refresh the same wait budget as its neighboring pause tests on slow Windows runners. Mirror production sidebar menu state in the browser fixture so opening its actions suppresses the hover preview and does not intercept the menu click. * test(literature): allow slow Windows refresh completion The shared refresh helper waits for a completed background run after scheduling classification. Under full Windows load the run can still be active at five seconds; use a bounded fifteen-second wait without changing application behavior. * fix(acp): order continuation admission before project locks Acquire root-session admission before the Project gate for parent messages, interrupted-turn continuation, and Save as skill. Opposite lock ordering could deadlock a queued root prompt against an upward child message. Keep Project deletion protection through validation, resume and provider acceptance, then release it without holding it through the whole turn. Cover the deterministic deadlock, error cleanup and completion fallback, and capture the successful concurrent-message Electron journey. * test(delegation): drain the verification owner before cleanup * test(browser): align preview layout and await menu focus
This commit is contained in:
@@ -163,6 +163,7 @@ useComputeStore.setState({
|
||||
|
||||
export function Fixture(): React.JSX.Element {
|
||||
const [actions, setActions] = useState<string[]>([])
|
||||
const [openSessionActionsId, setOpenSessionActionsId] = useState<string | null>(null)
|
||||
const record = (action: string): void => setActions((current) => [...current, action])
|
||||
useEffect(() => {
|
||||
useNavigationStore.setState({
|
||||
@@ -243,6 +244,12 @@ export function Fixture(): React.JSX.Element {
|
||||
onNewConversation={noop}
|
||||
onOpenFiles={noop}
|
||||
onOpenSession={() => record('session')}
|
||||
openSessionActionsId={openSessionActionsId}
|
||||
onSessionActionsOpenChange={(sessionId, open) => {
|
||||
setOpenSessionActionsId((current) =>
|
||||
open ? sessionId : current === sessionId ? null : current
|
||||
)
|
||||
}}
|
||||
onRenameSession={noop}
|
||||
onDownloadArtifacts={noop}
|
||||
onViewNotebook={noop}
|
||||
|
||||
@@ -22,7 +22,8 @@ for (const width of [320, 375, 414, 768]) {
|
||||
const literatureBounds = await literature.boundingBox()
|
||||
expect(expandBounds).not.toBeNull()
|
||||
expect(literatureBounds).not.toBeNull()
|
||||
expect(Math.abs(expandBounds!.y - literatureBounds!.y)).toBeLessThanOrEqual(1)
|
||||
// Literature navigation is a separate footer below the abstract expansion control.
|
||||
expect(literatureBounds!.y).toBeGreaterThanOrEqual(expandBounds!.y + expandBounds!.height)
|
||||
await expand.click()
|
||||
await expect(expand).toHaveAttribute('aria-expanded', 'true')
|
||||
expect(await page.evaluate(() => document.documentElement.scrollWidth <= innerWidth)).toBe(
|
||||
|
||||
@@ -33,12 +33,18 @@ for (const dark of [false, true]) {
|
||||
const metadata = (await row.getByText('1 session', { exact: true }).boundingBox())!
|
||||
await page.mouse.click(metadata.x + metadata.width / 2, metadata.y + metadata.height / 2)
|
||||
await expect(page.getByTestId('actions')).toHaveText('project,project,project,project,project')
|
||||
await page.getByRole('button', { name: 'Open actions for P1' }).click()
|
||||
const actions = page.getByRole('button', { name: 'Open actions for P1' })
|
||||
await actions.click()
|
||||
await expect(page.getByRole('menu')).toBeVisible()
|
||||
await page.keyboard.press('Escape')
|
||||
await expect(page.getByRole('menu')).toBeHidden()
|
||||
await expect(actions).toBeFocused()
|
||||
await expect(page.getByTestId('actions')).toHaveText('project,project,project,project,project')
|
||||
await project.focus()
|
||||
await expect(project).toBeFocused()
|
||||
await page.keyboard.press('Enter')
|
||||
await expect(page.getByTestId('actions')).toHaveText(Array(6).fill('project').join(','))
|
||||
await expect(project).toBeFocused()
|
||||
await page.keyboard.press('Space')
|
||||
await expect(page.getByTestId('actions')).toHaveText(Array(7).fill('project').join(','))
|
||||
// The decorative card inset is outside any row and must remain inert.
|
||||
@@ -55,6 +61,7 @@ test('sidebar edges open the session while the dropdown and context menu remain
|
||||
const row = page.locator('[data-session-id="session-1"]')
|
||||
await clickEdges(page, row, 'session')
|
||||
await page.getByRole('button', { name: 'Open actions for Analysis session' }).click()
|
||||
await expect(page.getByTestId('session-preview-content')).toBeHidden()
|
||||
await page.getByRole('menuitem', { name: /Pin/ }).click()
|
||||
await expect(page.getByTestId('actions')).toHaveText('session,session,session,session,pin')
|
||||
await row.click({ button: 'right', position: { x: 3, y: 2 } })
|
||||
|
||||
@@ -1000,7 +1000,9 @@ test('recovers a post-fence receipt persistence failure as uncertain after proce
|
||||
expect(receipt).toMatchObject({ status: 'uncertain', resolution: 'pending' })
|
||||
})
|
||||
|
||||
test('fairly schedules two upward lanes with a concurrent real user prompt', async ({ app }) => {
|
||||
test('fairly schedules two upward lanes with a concurrent real user prompt', async ({
|
||||
app
|
||||
}, testInfo) => {
|
||||
test.setTimeout(180_000)
|
||||
await app.completeOnboarding()
|
||||
const page = await app.configureFakeAgent()
|
||||
@@ -1056,6 +1058,7 @@ test('fairly schedules two upward lanes with a concurrent real user prompt', asy
|
||||
{ requestId: 'e2e-fairness-a', status: 'accepted' },
|
||||
{ requestId: 'e2e-fairness-b', status: 'accepted' }
|
||||
])
|
||||
await page.screenshot({ path: testInfo.outputPath('concurrent-parent-messages.png') })
|
||||
})
|
||||
|
||||
test('stops only the active branch and exposes a retryable partial failure', async ({
|
||||
|
||||
Generated
+14
-2
@@ -11137,11 +11137,23 @@
|
||||
}
|
||||
},
|
||||
"node_modules/deepmerge-ts": {
|
||||
"version": "7.1.5",
|
||||
"version": "8.0.2",
|
||||
"resolved": "https://registry.npmjs.org/deepmerge-ts/-/deepmerge-ts-8.0.2.tgz",
|
||||
"integrity": "sha512-uqbvqLUMrc6p0MO+WBRtTxY55hmyh94WRwI5a++PZe54X+bfVh59FSN7uWCBCW1CCVjzjnrwzfI8zidE2obMMw==",
|
||||
"devOptional": true,
|
||||
"funding": [
|
||||
{
|
||||
"type": "ko-fi",
|
||||
"url": "https://ko-fi.com/rebeccastevens"
|
||||
},
|
||||
{
|
||||
"type": "tidelift",
|
||||
"url": "https://tidelift.com/funding/github/npm/deepmerge-ts"
|
||||
}
|
||||
],
|
||||
"license": "BSD-3-Clause",
|
||||
"engines": {
|
||||
"node": ">=16.0.0"
|
||||
"node": ">=16.9.0"
|
||||
}
|
||||
},
|
||||
"node_modules/default-browser": {
|
||||
|
||||
@@ -106,6 +106,11 @@
|
||||
"ws": "^8.21.1",
|
||||
"zod": "^4.4.3"
|
||||
},
|
||||
"overrides": {
|
||||
"@prisma/config@6.19.3": {
|
||||
"deepmerge-ts": "8.0.2"
|
||||
}
|
||||
},
|
||||
"devDependencies": {
|
||||
"@aiden0z/pptx-renderer": "1.2.4",
|
||||
"@anthropic-ai/claude-agent-sdk": "0.3.232",
|
||||
|
||||
@@ -0,0 +1,113 @@
|
||||
import { spawnSync } from 'node:child_process'
|
||||
import { mkdtempSync, rmSync, writeFileSync } from 'node:fs'
|
||||
import { createRequire } from 'node:module'
|
||||
import { tmpdir } from 'node:os'
|
||||
import { join } from 'node:path'
|
||||
|
||||
import { afterEach, describe, expect, it } from 'vitest'
|
||||
|
||||
// Resolve through the real consumer, so a vulnerable nested copy cannot be hidden
|
||||
// by a fixed root dependency. Native Node also exercises Prisma's dynamic imports.
|
||||
const require = createRequire(import.meta.url)
|
||||
const prismaRequire = createRequire(require.resolve('prisma/package.json'))
|
||||
const configPath = prismaRequire.resolve('@prisma/config')
|
||||
const prelude = `
|
||||
const assert = require('node:assert/strict');
|
||||
const { createRequire } = require('node:module');
|
||||
const configRequire = createRequire(${JSON.stringify(configPath)});
|
||||
const { deepmerge, deepmergeInto } = configRequire('deepmerge-ts');
|
||||
`
|
||||
let fixture: string | undefined
|
||||
|
||||
function runNode(source: string): void {
|
||||
const result = spawnSync(process.execPath, ['-e', prelude + source], {
|
||||
encoding: 'utf8',
|
||||
timeout: 10_000
|
||||
})
|
||||
expect(result.error, result.stderr).toBeUndefined()
|
||||
expect(result.status, result.stderr || result.stdout).toBe(0)
|
||||
}
|
||||
|
||||
afterEach(() => {
|
||||
if (fixture) rmSync(fixture, { recursive: true, force: true })
|
||||
fixture = undefined
|
||||
})
|
||||
|
||||
describe('Prisma config dependency compatibility', () => {
|
||||
it('handles recursive records without stack exhaustion (CVE-2026-40345)', () => {
|
||||
runNode(`
|
||||
const left = { left: true };
|
||||
left.self = left;
|
||||
const right = { right: true };
|
||||
right.self = right;
|
||||
const merged = deepmerge(left, right);
|
||||
assert.equal(merged.left, true);
|
||||
assert.equal(merged.right, true);
|
||||
assert.equal(merged.self, merged);
|
||||
const target = {};
|
||||
target.self = target;
|
||||
deepmergeInto(target, right);
|
||||
assert.equal(target.right, true);
|
||||
assert.equal(target.self, target);
|
||||
`)
|
||||
})
|
||||
|
||||
it('preserves plain config merge values and leaves inputs unchanged', () => {
|
||||
runNode(`
|
||||
const left = {
|
||||
schema: 'old.prisma',
|
||||
migrations: { path: 'migrations', seed: 'node seed.mjs' },
|
||||
tables: { external: ['existing'] }
|
||||
};
|
||||
const right = {
|
||||
schema: 'schema.prisma',
|
||||
migrations: { seed: 'node next.mjs' },
|
||||
tables: { external: ['next'] }
|
||||
};
|
||||
const before = JSON.stringify([left, right]);
|
||||
assert.deepEqual(deepmerge(left, right), {
|
||||
schema: 'schema.prisma',
|
||||
migrations: { path: 'migrations', seed: 'node next.mjs' },
|
||||
tables: { external: ['existing', 'next'] }
|
||||
});
|
||||
assert.equal(JSON.stringify([left, right]), before);
|
||||
`)
|
||||
})
|
||||
|
||||
it.each(['cjs', 'mts'])(
|
||||
'loads a real %s config and preserves missing-config behavior',
|
||||
(extension) => {
|
||||
fixture = mkdtempSync(join(tmpdir(), 'prisma-config-compatibility-'))
|
||||
const configFile = join(fixture, `prisma.config.${extension}`)
|
||||
writeFileSync(
|
||||
configFile,
|
||||
`${extension === 'cjs' ? 'module.exports =' : 'export default'} {
|
||||
schema: './prisma/schema.prisma',
|
||||
migrations: { path: './prisma/migrations', seed: 'node seed.mjs' }
|
||||
}`
|
||||
)
|
||||
runNode(`
|
||||
const { join } = require('node:path');
|
||||
const { loadConfigFromFile } = require(${JSON.stringify(configPath)});
|
||||
const root = ${JSON.stringify(fixture)};
|
||||
(async () => {
|
||||
const loaded = await loadConfigFromFile({
|
||||
configRoot: root, configFile: ${JSON.stringify(configFile)}
|
||||
});
|
||||
assert.equal(loaded.error, undefined);
|
||||
assert.equal(loaded.resolvedPath, ${JSON.stringify(configFile)});
|
||||
assert.equal(loaded.config.schema, join(root, 'prisma', 'schema.prisma'));
|
||||
assert.equal(loaded.config.migrations.path, join(root, 'prisma', 'migrations'));
|
||||
assert.equal(loaded.config.migrations.seed, 'node seed.mjs');
|
||||
const absent = await loadConfigFromFile({ configRoot: join(root, 'absent') });
|
||||
assert.equal(absent.error, undefined);
|
||||
assert.equal(absent.resolvedPath, null);
|
||||
const missing = await loadConfigFromFile({
|
||||
configRoot: root, configFile: join(root, 'missing.config.ts')
|
||||
});
|
||||
assert.equal(missing.error._tag, 'ConfigFileNotFound');
|
||||
})().catch((error) => { console.error(error); process.exitCode = 1; });
|
||||
`)
|
||||
}
|
||||
)
|
||||
})
|
||||
@@ -126,9 +126,18 @@ const createHarness = (
|
||||
)
|
||||
const snapshot = { status: 'connected' } as never
|
||||
const startContinuationWhenDispatchAdmitted = vi.fn(
|
||||
async (request: unknown, validate: () => Promise<void>) => {
|
||||
await validate()
|
||||
return startContinuation(request)
|
||||
async (
|
||||
request: unknown,
|
||||
validate: () => Promise<void>,
|
||||
_messageId?: string,
|
||||
_queued?: () => void,
|
||||
admitDispatch?: (operation: () => Promise<void>) => Promise<void>
|
||||
) => {
|
||||
const dispatch = async (): Promise<void> => {
|
||||
await validate()
|
||||
await startContinuation(request)
|
||||
}
|
||||
return admitDispatch ? admitDispatch(dispatch) : dispatch()
|
||||
}
|
||||
)
|
||||
const workflows = createAcpHandlerWorkflows(
|
||||
@@ -255,6 +264,55 @@ describe('ACP resume Session workflow', () => {
|
||||
})
|
||||
})
|
||||
|
||||
describe('continuation admission ordering', () => {
|
||||
it.each(['continue', 'save-as-skill'] as const)(
|
||||
'queues %s at the root before holding Project admission',
|
||||
async (kind) => {
|
||||
let held = false
|
||||
const harness = createHarness(undefined, {
|
||||
withSessionAvailable: async (_projectId, _sessionId, operation) => {
|
||||
expect(held).toBe(false)
|
||||
held = true
|
||||
try {
|
||||
return await operation()
|
||||
} finally {
|
||||
held = false
|
||||
}
|
||||
},
|
||||
withSessionAvailableById: vi.fn()
|
||||
})
|
||||
if (kind === 'continue') {
|
||||
harness.session.status = 'error'
|
||||
harness.session.activeRun = undefined
|
||||
harness.session.resumeRecovery = {
|
||||
kind: 'resume-required',
|
||||
cause: 'app-restart',
|
||||
promptMessageId: 'prompt-1'
|
||||
}
|
||||
}
|
||||
harness.startContinuationWhenDispatchAdmitted.mockImplementationOnce(
|
||||
async (_request, validate, _messageId, _queued, admitDispatch) => {
|
||||
expect(held).toBe(false)
|
||||
expect(admitDispatch).toEqual(expect.any(Function))
|
||||
await admitDispatch!(async () => {
|
||||
expect(held).toBe(true)
|
||||
await validate()
|
||||
})
|
||||
}
|
||||
)
|
||||
if (kind === 'continue')
|
||||
await harness.workflows.continueInterruptedTurn({
|
||||
projectId: 'project-1',
|
||||
sessionId: 'session-1',
|
||||
promptMessageId: 'prompt-1'
|
||||
})
|
||||
else await harness.workflows.saveAsSkill(harness.request)
|
||||
expect(harness.startContinuationWhenDispatchAdmitted).toHaveBeenCalledOnce()
|
||||
expect(held).toBe(false)
|
||||
}
|
||||
)
|
||||
})
|
||||
|
||||
describe('ACP interrupted turn workflow', () => {
|
||||
it('does not re-enter archive admission while continuing an interrupted turn', async () => {
|
||||
let archiveQueue: Promise<void> = Promise.resolve()
|
||||
@@ -283,9 +341,8 @@ describe('ACP interrupted turn workflow', () => {
|
||||
// itself; the dispatch-admitted path intentionally bypasses only that nested guard.
|
||||
harness.startContinuation.mockImplementationOnce(() => enqueueArchive(async () => undefined))
|
||||
harness.startContinuationWhenDispatchAdmitted.mockImplementationOnce(
|
||||
async (_request: unknown, validate: () => Promise<void>) => {
|
||||
await validate()
|
||||
return 'provider_prompt_accepted'
|
||||
async (_request, validate, _messageId, _queued, admitDispatch) => {
|
||||
await admitDispatch!(validate)
|
||||
}
|
||||
)
|
||||
|
||||
@@ -499,7 +556,7 @@ describe('ACP Save as skill workflow', () => {
|
||||
}
|
||||
)
|
||||
|
||||
it('dispatches through the Session admission already held by the workflow', async () => {
|
||||
it('passes Session admission into the root continuation scheduler', async () => {
|
||||
const harness = createHarness()
|
||||
|
||||
await harness.workflows.saveAsSkill(harness.request)
|
||||
|
||||
@@ -49,7 +49,10 @@ type AcpHandlerWorkflowRuntime = {
|
||||
startContinuation(request: AcpPromptRequest): Promise<void>
|
||||
startContinuationWhenDispatchAdmitted(
|
||||
request: AcpPromptRequest,
|
||||
validate: () => Promise<void>
|
||||
validate: () => Promise<void>,
|
||||
delegatedMessageId?: string,
|
||||
onAdmissionQueued?: () => void,
|
||||
admitDispatch?: (operation: () => Promise<void>) => Promise<void>
|
||||
): Promise<unknown>
|
||||
}
|
||||
|
||||
@@ -345,16 +348,32 @@ const createAcpHandlerWorkflows = (
|
||||
startDispatchAdmittedContinuation: (
|
||||
continuation: AcpPromptRequest,
|
||||
validate: () => Promise<void>
|
||||
) => runtime.startContinuationWhenDispatchAdmitted(continuation, validate)
|
||||
) =>
|
||||
runtime.startContinuationWhenDispatchAdmitted(
|
||||
continuation,
|
||||
validate,
|
||||
undefined,
|
||||
undefined,
|
||||
(operation) =>
|
||||
archiveAvailability.withSessionAvailable(
|
||||
request.projectId,
|
||||
request.sessionId,
|
||||
operation
|
||||
)
|
||||
)
|
||||
}
|
||||
: {}),
|
||||
notifications: taskNotifications
|
||||
},
|
||||
request
|
||||
)
|
||||
return archiveAvailability
|
||||
? archiveAvailability.withSessionAvailable(request.projectId, request.sessionId, run)
|
||||
: run()
|
||||
// Check availability before reading, but never hold the Project gate while queuing root work.
|
||||
await archiveAvailability?.withSessionAvailable(
|
||||
request.projectId,
|
||||
request.sessionId,
|
||||
async () => undefined
|
||||
)
|
||||
return run()
|
||||
},
|
||||
|
||||
async saveAsSkill(request): Promise<AcpRuntimeState> {
|
||||
@@ -370,26 +389,42 @@ const createAcpHandlerWorkflows = (
|
||||
text: 'Save as skill'
|
||||
})
|
||||
try {
|
||||
await runtime.startContinuationWhenDispatchAdmitted(prepared.continuation, async () => {
|
||||
await saveAsSkillAdmission?.(request.sessionId)
|
||||
const admitted = prepareSaveAsSkillContinuation(
|
||||
runtime,
|
||||
await interruptedTurnSessions.loadSession(request.projectId, request.sessionId),
|
||||
request
|
||||
)
|
||||
if (!isDeepStrictEqual(admitted.continuation, prepared.continuation)) {
|
||||
throw new Error('Save as skill Session changed before provider admission.')
|
||||
}
|
||||
})
|
||||
await runtime.startContinuationWhenDispatchAdmitted(
|
||||
prepared.continuation,
|
||||
async () => {
|
||||
await saveAsSkillAdmission?.(request.sessionId)
|
||||
const admitted = prepareSaveAsSkillContinuation(
|
||||
runtime,
|
||||
await interruptedTurnSessions.loadSession(request.projectId, request.sessionId),
|
||||
request
|
||||
)
|
||||
if (!isDeepStrictEqual(admitted.continuation, prepared.continuation)) {
|
||||
throw new Error('Save as skill Session changed before provider admission.')
|
||||
}
|
||||
},
|
||||
undefined,
|
||||
undefined,
|
||||
archiveAvailability
|
||||
? (operation) =>
|
||||
archiveAvailability.withSessionAvailable(
|
||||
request.projectId,
|
||||
request.sessionId,
|
||||
operation
|
||||
)
|
||||
: undefined
|
||||
)
|
||||
} catch (error) {
|
||||
if (tracked) taskNotifications?.untrackPrompt(prepared.session.id, tracked)
|
||||
throw error
|
||||
}
|
||||
return runtime.getState()
|
||||
}
|
||||
return archiveAvailability
|
||||
? archiveAvailability.withSessionAvailable(request.projectId, request.sessionId, save)
|
||||
: save()
|
||||
await archiveAvailability?.withSessionAvailable(
|
||||
request.projectId,
|
||||
request.sessionId,
|
||||
async () => undefined
|
||||
)
|
||||
return save()
|
||||
},
|
||||
|
||||
async sendPrompt(request): Promise<AcpRuntimeState> {
|
||||
|
||||
@@ -3156,6 +3156,137 @@ describe('AcpRuntimeCoordinator', () => {
|
||||
expect(created.sendPrompt).toHaveBeenCalledOnce()
|
||||
})
|
||||
|
||||
it.each(['claude-code', 'opencode', 'codex'] as const)(
|
||||
'admits an upward message without reversing the root and Project lock order for %s',
|
||||
async (frameworkId) => {
|
||||
let created!: ReturnType<typeof createFakeRuntime>
|
||||
const userFinished = createDeferred<unknown>()
|
||||
const parentFinished = createDeferred<unknown>()
|
||||
const parentAtProvider = createDeferred()
|
||||
const allowParentAcceptance = createDeferred()
|
||||
let promptIndex = 0
|
||||
const coordinator = new AcpRuntimeCoordinator((callbacks) => {
|
||||
created = createFakeRuntime({
|
||||
frameworkId,
|
||||
sessionIds: ['session-1'],
|
||||
callbacks,
|
||||
beforeProviderPromptAccepted: async () => {
|
||||
if (promptIndex === 2) {
|
||||
parentAtProvider.resolve()
|
||||
await allowParentAcceptance.promise
|
||||
}
|
||||
},
|
||||
prompt: () => (promptIndex++ === 0 ? userFinished.promise : parentFinished.promise)
|
||||
})
|
||||
return created.runtime
|
||||
})
|
||||
const session = await coordinator.createSession({ cwd: '/workspace', projectId: 'project-1' })
|
||||
const archive = new ArchiveCoordinator(
|
||||
{ get: vi.fn(), updateArchive: vi.fn() },
|
||||
{
|
||||
sessionProjectId: async () => 'project-1',
|
||||
assertProjectArchivable: vi.fn(),
|
||||
assertSessionAvailable: vi.fn(),
|
||||
updateArchive: vi.fn()
|
||||
},
|
||||
{
|
||||
isSessionBusy: () => false,
|
||||
isProjectBusy: () => false,
|
||||
liveSessionProjectId: () => 'project-1'
|
||||
}
|
||||
)
|
||||
const userAtDispatch = createDeferred()
|
||||
const allowUserDispatch = createDeferred()
|
||||
coordinator.setPromptDispatchAdmissionGuard(async (sessionId, dispatch) => {
|
||||
userAtDispatch.resolve()
|
||||
await allowUserDispatch.promise
|
||||
return archive.withSessionDeletionAdmissionById(sessionId, dispatch)
|
||||
})
|
||||
const user = coordinator.sendPrompt({ sessionId: session.sessionId, text: 'concurrent user' })
|
||||
await userAtDispatch.promise
|
||||
const parentAtProject = createDeferred()
|
||||
const parentQueued = createDeferred()
|
||||
const upward = coordinator.startContinuationWhenDispatchAdmitted(
|
||||
{ sessionId: session.sessionId, text: 'upward message' },
|
||||
async () => undefined,
|
||||
'message-1',
|
||||
() => parentQueued.resolve(),
|
||||
(operation) =>
|
||||
archive.withProjectDeletionAdmission('project-1', async () => {
|
||||
parentAtProject.resolve()
|
||||
await operation()
|
||||
})
|
||||
)
|
||||
await parentQueued.promise
|
||||
allowUserDispatch.resolve()
|
||||
await vi.waitFor(() => expect(created.sendPrompt).toHaveBeenCalledOnce())
|
||||
userFinished.resolve({ stopReason: 'end_turn' })
|
||||
await user
|
||||
await parentAtProject.promise
|
||||
await parentAtProvider.promise
|
||||
const nextProjectOperation = vi.fn(async () => 'available')
|
||||
const projectAvailable = archive.withProjectDeletionAdmission(
|
||||
'project-1',
|
||||
nextProjectOperation
|
||||
)
|
||||
await Promise.resolve()
|
||||
expect(nextProjectOperation).not.toHaveBeenCalled()
|
||||
allowParentAcceptance.resolve()
|
||||
await expect(upward).resolves.toBe('provider_prompt_accepted')
|
||||
// Acceptance releases the Project gate even while the provider turn is still running.
|
||||
await expect(projectAvailable).resolves.toBe('available')
|
||||
const laterUser = coordinator.sendPrompt({ sessionId: session.sessionId, text: 'later user' })
|
||||
await Promise.resolve()
|
||||
expect(created.sendPrompt).toHaveBeenCalledOnce()
|
||||
parentFinished.resolve({ stopReason: 'end_turn' })
|
||||
await laterUser
|
||||
expect(created.sendPrompt).toHaveBeenCalledTimes(2)
|
||||
}
|
||||
)
|
||||
|
||||
it.each(['admission', 'validation', 'provider'] as const)(
|
||||
'releases parent-message deletion admission after a %s failure',
|
||||
async (failurePhase) => {
|
||||
const failure = new Error('parent delivery failed')
|
||||
const coordinator = new AcpRuntimeCoordinator(
|
||||
(callbacks) =>
|
||||
createFakeRuntime({
|
||||
frameworkId: 'opencode',
|
||||
sessionIds: ['session-1'],
|
||||
callbacks,
|
||||
beforeProviderPromptAccepted: async () => {
|
||||
if (failurePhase === 'provider') throw failure
|
||||
}
|
||||
}).runtime
|
||||
)
|
||||
const session = await coordinator.createSession({ cwd: '/workspace' })
|
||||
let held = false
|
||||
const released = createDeferred()
|
||||
const upward = coordinator.startContinuationWhenDispatchAdmitted(
|
||||
{ sessionId: session.sessionId, text: 'upward message' },
|
||||
async () => {
|
||||
expect(held).toBe(true)
|
||||
if (failurePhase === 'validation') throw failure
|
||||
},
|
||||
'message-1',
|
||||
undefined,
|
||||
async (operation) => {
|
||||
held = true
|
||||
try {
|
||||
if (failurePhase === 'admission') throw failure
|
||||
await operation()
|
||||
} finally {
|
||||
held = false
|
||||
released.resolve()
|
||||
}
|
||||
}
|
||||
)
|
||||
await expect(upward).rejects.toThrow('parent delivery failed')
|
||||
await released.promise
|
||||
expect(held).toBe(false)
|
||||
}
|
||||
)
|
||||
|
||||
it('linearizes real user prompts and upward continuations through one root admission lock', async () => {
|
||||
const prompts = [
|
||||
createDeferred<unknown>(),
|
||||
@@ -3288,25 +3419,36 @@ describe('AcpRuntimeCoordinator', () => {
|
||||
expect(created.sendAppContinuation).not.toHaveBeenCalled()
|
||||
})
|
||||
|
||||
it('returns completion evidence when an upward continuation completes without an acceptance callback', async () => {
|
||||
const coordinator = new AcpRuntimeCoordinator(
|
||||
(callbacks) =>
|
||||
createFakeRuntime({
|
||||
frameworkId: 'codex',
|
||||
sessionIds: ['session-1'],
|
||||
callbacks,
|
||||
skipProviderPromptAccepted: true
|
||||
}).runtime
|
||||
)
|
||||
const session = await coordinator.createSession({ cwd: '/workspace' })
|
||||
|
||||
await expect(
|
||||
coordinator.startContinuationWhen(
|
||||
{ sessionId: session.sessionId, text: 'completion fallback' },
|
||||
async () => undefined
|
||||
it.each([false, true])(
|
||||
'returns completion evidence without an acceptance callback (guarded=%s)',
|
||||
async (guarded) => {
|
||||
const coordinator = new AcpRuntimeCoordinator(
|
||||
(callbacks) =>
|
||||
createFakeRuntime({
|
||||
frameworkId: 'codex',
|
||||
sessionIds: ['session-1'],
|
||||
callbacks,
|
||||
skipProviderPromptAccepted: true
|
||||
}).runtime
|
||||
)
|
||||
).resolves.toBe('provider_prompt_completed')
|
||||
})
|
||||
const session = await coordinator.createSession({ cwd: '/workspace' })
|
||||
|
||||
await expect(
|
||||
guarded
|
||||
? coordinator.startContinuationWhenDispatchAdmitted(
|
||||
{ sessionId: session.sessionId, text: 'completion fallback' },
|
||||
async () => undefined,
|
||||
undefined,
|
||||
undefined,
|
||||
(operation) => operation()
|
||||
)
|
||||
: coordinator.startContinuationWhen(
|
||||
{ sessionId: session.sessionId, text: 'completion fallback' },
|
||||
async () => undefined
|
||||
)
|
||||
).resolves.toBe('provider_prompt_completed')
|
||||
}
|
||||
)
|
||||
|
||||
it('rejects a user prompt still waiting on admission when quit begins', async () => {
|
||||
const admission = createDeferred<void>()
|
||||
|
||||
@@ -1062,16 +1062,19 @@ class AcpRuntimeCoordinator {
|
||||
request: AcpPromptRequest,
|
||||
validate: () => Promise<void>,
|
||||
delegatedMessageId?: string,
|
||||
onAdmissionQueued?: () => void
|
||||
onAdmissionQueued?: () => void,
|
||||
admitDispatch?: (operation: () => Promise<void>) => Promise<void>
|
||||
): Promise<DelegateMessageAcceptanceEvidence> {
|
||||
// The caller owns final deletion admission for the whole validation/resume/acceptance lifecycle.
|
||||
// Bypass only the nested dispatch guard; root-session admission remains linearized below.
|
||||
// Acquire deletion admission inside root admission, retaining it through validation, resume,
|
||||
// and acceptance. Taking the Project gate first can deadlock a user prompt waiting for it.
|
||||
// Legacy already-admitted callers can omit the wrapper; nested dispatch admission is bypassed.
|
||||
return this.startContinuationWhenWithDispatchAdmission(
|
||||
request,
|
||||
validate,
|
||||
true,
|
||||
delegatedMessageId,
|
||||
onAdmissionQueued
|
||||
onAdmissionQueued,
|
||||
admitDispatch
|
||||
)
|
||||
}
|
||||
|
||||
@@ -1080,7 +1083,8 @@ class AcpRuntimeCoordinator {
|
||||
validate: () => Promise<void>,
|
||||
dispatchAdmitted: boolean,
|
||||
delegatedMessageId?: string,
|
||||
onAdmissionQueued?: () => void
|
||||
onAdmissionQueued?: () => void,
|
||||
admitDispatch?: (operation: () => Promise<void>) => Promise<void>
|
||||
): Promise<DelegateMessageAcceptanceEvidence> {
|
||||
let resolve!: (evidence: DelegateMessageAcceptanceEvidence) => void
|
||||
let reject!: (error: unknown) => void
|
||||
@@ -1106,7 +1110,7 @@ class AcpRuntimeCoordinator {
|
||||
reject(error)
|
||||
}
|
||||
|
||||
const admission = this.linearizeRootAdmission(request.sessionId, async () => {
|
||||
const dispatch = async (): Promise<void> => {
|
||||
try {
|
||||
await validate()
|
||||
} catch (error) {
|
||||
@@ -1135,6 +1139,17 @@ class AcpRuntimeCoordinator {
|
||||
acceptance.settled = true
|
||||
resolve('provider_prompt_completed')
|
||||
}
|
||||
}
|
||||
const admission = this.linearizeRootAdmission(request.sessionId, async () => {
|
||||
if (!admitDispatch) return dispatch()
|
||||
let completion!: Promise<void>
|
||||
await admitDispatch(async () => {
|
||||
completion = dispatch()
|
||||
void completion.catch((error) => acceptance.reject(error))
|
||||
// Release the Project gate at acceptance, while root admission still owns the whole turn.
|
||||
await accepted
|
||||
})
|
||||
await completion
|
||||
})
|
||||
onAdmissionQueued?.()
|
||||
void admission.catch((error) => acceptance.reject(error))
|
||||
|
||||
@@ -1664,7 +1664,9 @@ describe('production delegated-work composition', () => {
|
||||
|
||||
expect(first).toMatchObject({ kind: 'receipts', children: [{ name: 'Source audit' }] })
|
||||
expect(second).toMatchObject({ kind: 'receipts', children: [{ name: 'Source audit 2' }] })
|
||||
await expect(harness.reopen().host.children(harness.caller)).resolves.toMatchObject([
|
||||
const verification = harness.reopen()
|
||||
compositions.push(verification)
|
||||
await expect(verification.host.children(harness.caller)).resolves.toMatchObject([
|
||||
{ name: 'Source audit' },
|
||||
{ name: 'Source audit 2' }
|
||||
])
|
||||
@@ -1675,7 +1677,7 @@ describe('production delegated-work composition', () => {
|
||||
.map(({ delegateName }) => delegateName)
|
||||
).toEqual(['Source audit', 'Source audit 2'])
|
||||
} finally {
|
||||
// Both owners can still be materializing the wait:false children. Drain them before rm.
|
||||
// Reopened owners can still be recovering the wait:false children. Drain them before rm.
|
||||
for (const composition of compositions) await composition.root.shutdown()
|
||||
}
|
||||
})
|
||||
|
||||
+95
-102
@@ -2996,117 +2996,110 @@ const createApplicationModules = async (
|
||||
},
|
||||
parentMessages: {
|
||||
async deliver(delivery) {
|
||||
return archiveCoordinator.withProjectDeletionAdmission(
|
||||
const runtime = runtimeRef.current
|
||||
if (!runtime) throw new Error('ACP runtime is not available.')
|
||||
const session = await sessionRepository.loadSession(
|
||||
delivery.session.projectId,
|
||||
delivery.session.sessionId
|
||||
)
|
||||
const graph = session?.conversationGraph
|
||||
const rootFrame = graph?.frames.find((frame) => frame.id === delivery.targetFrameId)
|
||||
const rootBranch = graph?.branches.find((branch) => branch.id === rootFrame?.activeBranchId)
|
||||
if (
|
||||
!session ||
|
||||
session.id !== delivery.session.sessionId ||
|
||||
session.projectId !== delivery.session.projectId ||
|
||||
graph?.rootFrameId !== delivery.targetFrameId ||
|
||||
!rootBranch ||
|
||||
!graph.messages.some((message) => message.id === delivery.originMessageId)
|
||||
) {
|
||||
throw new Error('Parent message durable root provenance is unavailable.')
|
||||
}
|
||||
return runtime.startContinuationWhenDispatchAdmitted(
|
||||
{
|
||||
sessionId: delivery.session.sessionId,
|
||||
text:
|
||||
`[Delegated ${delivery.kind} from Frame ${delivery.sourceFrameId}, ` +
|
||||
`Attempt ${delivery.sourceAttemptId}]\n\n${delivery.text}`,
|
||||
suppressUserMessage: true,
|
||||
provenanceContext: {
|
||||
// Suppressed continuations create no user node; replies retain the durable origin.
|
||||
promptMessageId: delivery.originMessageId,
|
||||
originMessageId: delivery.originMessageId,
|
||||
rootFrameId: graph.rootFrameId,
|
||||
agentFrameId: graph.rootFrameId,
|
||||
messageBranchId: delivery.rootBranchId,
|
||||
messageBranchAncestry: [delivery.rootBranchId],
|
||||
messageAncestry: [delivery.originMessageId],
|
||||
runtimeSegmentId: `delegated-message-${delivery.messageId}`
|
||||
}
|
||||
},
|
||||
async () => {
|
||||
const runtime = runtimeRef.current
|
||||
if (!runtime) throw new Error('ACP runtime is not available.')
|
||||
const session = await sessionRepository.loadSession(
|
||||
let latest = await sessionRepository.loadSession(
|
||||
delivery.session.projectId,
|
||||
delivery.session.sessionId
|
||||
)
|
||||
const graph = session?.conversationGraph
|
||||
const rootFrame = graph?.frames.find((frame) => frame.id === delivery.targetFrameId)
|
||||
const rootBranch = graph?.branches.find(
|
||||
(branch) => branch.id === rootFrame?.activeBranchId
|
||||
const latestGraph = latest?.conversationGraph
|
||||
const latestRoot = latestGraph?.frames.find(({ id }) => id === delivery.targetFrameId)
|
||||
const latestBranch = latestGraph?.branches.find(
|
||||
({ id }) => id === latestRoot?.activeBranchId
|
||||
)
|
||||
if (
|
||||
!session ||
|
||||
session.id !== delivery.session.sessionId ||
|
||||
session.projectId !== delivery.session.projectId ||
|
||||
graph?.rootFrameId !== delivery.targetFrameId ||
|
||||
!rootBranch ||
|
||||
!graph.messages.some((message) => message.id === delivery.originMessageId)
|
||||
!latest ||
|
||||
latestBranch?.id !== delivery.rootBranchId ||
|
||||
`${latestBranch.id}:${latestBranch.createdAt}` !== delivery.rootBranchRevision
|
||||
) {
|
||||
throw new Error('Parent message durable root provenance is unavailable.')
|
||||
throw new DelegateMessageParkedError(
|
||||
'Parent message root Branch changed before dispatch.'
|
||||
)
|
||||
}
|
||||
return runtime.startContinuationWhenDispatchAdmitted(
|
||||
{
|
||||
sessionId: delivery.session.sessionId,
|
||||
text:
|
||||
`[Delegated ${delivery.kind} from Frame ${delivery.sourceFrameId}, ` +
|
||||
`Attempt ${delivery.sourceAttemptId}]\n\n${delivery.text}`,
|
||||
suppressUserMessage: true,
|
||||
provenanceContext: {
|
||||
// Suppressed continuations create no user node; replies retain the durable origin.
|
||||
promptMessageId: delivery.originMessageId,
|
||||
originMessageId: delivery.originMessageId,
|
||||
rootFrameId: graph.rootFrameId,
|
||||
agentFrameId: graph.rootFrameId,
|
||||
messageBranchId: delivery.rootBranchId,
|
||||
messageBranchAncestry: [delivery.rootBranchId],
|
||||
messageAncestry: [delivery.originMessageId],
|
||||
runtimeSegmentId: `delegated-message-${delivery.messageId}`
|
||||
}
|
||||
},
|
||||
async () => {
|
||||
let latest = await sessionRepository.loadSession(
|
||||
delivery.session.projectId,
|
||||
delivery.session.sessionId
|
||||
)
|
||||
const latestGraph = latest?.conversationGraph
|
||||
const latestRoot = latestGraph?.frames.find(
|
||||
({ id }) => id === delivery.targetFrameId
|
||||
)
|
||||
const latestBranch = latestGraph?.branches.find(
|
||||
({ id }) => id === latestRoot?.activeBranchId
|
||||
)
|
||||
if (
|
||||
!latest ||
|
||||
latestBranch?.id !== delivery.rootBranchId ||
|
||||
`${latestBranch.id}:${latestBranch.createdAt}` !== delivery.rootBranchRevision
|
||||
) {
|
||||
throw new DelegateMessageParkedError(
|
||||
'Parent message root Branch changed before dispatch.'
|
||||
)
|
||||
}
|
||||
const agentTarget = await resolveSessionAgentTarget(latest)
|
||||
if (
|
||||
agentTarget &&
|
||||
shouldPersistSessionAgentConfiguration(latest.agentConfiguration, agentTarget)
|
||||
) {
|
||||
latest = await sessionPersistenceCoordinator.saveSession({
|
||||
...latest,
|
||||
agentConfiguration: toSessionAgentConfiguration(agentTarget)
|
||||
})
|
||||
}
|
||||
if (!runtime.hasLiveSession(latest.projectId, latest.id) || agentTarget) {
|
||||
await runtime.resumeSession({
|
||||
sessionId: latest.id,
|
||||
cwd: latest.cwd,
|
||||
projectId: latest.projectId,
|
||||
...(latest.permissionProfile
|
||||
? { permissionProfile: latest.permissionProfile }
|
||||
: {}),
|
||||
memoryEnabled: latest.memoryEnabled !== false,
|
||||
...(latest.agentFrameworkId
|
||||
? { previousFrameworkId: latest.agentFrameworkId }
|
||||
: {}),
|
||||
...(latest.agentBackendId ? { previousBackendId: latest.agentBackendId } : {}),
|
||||
...(latest.specialistId ? { specialistId: latest.specialistId } : {}),
|
||||
...(latest.specialistBindingPending === true
|
||||
? { specialistBindingPending: true }
|
||||
: {}),
|
||||
...(latest.providerSessionId
|
||||
? { providerSessionId: latest.providerSessionId }
|
||||
: {}),
|
||||
...(latest.providerContinuityToken
|
||||
? { providerContinuityToken: latest.providerContinuityToken }
|
||||
: {}),
|
||||
...(agentTarget ? { agentTarget } : {})
|
||||
})
|
||||
}
|
||||
const started = await delivery.startDispatch()
|
||||
if (started !== 'started') {
|
||||
throw new DelegateMessageParkedError(
|
||||
'Parent message dispatch fence was not acquired.'
|
||||
)
|
||||
}
|
||||
},
|
||||
delivery.messageId,
|
||||
delivery.onRootAdmissionQueued
|
||||
)
|
||||
}
|
||||
const agentTarget = await resolveSessionAgentTarget(latest)
|
||||
if (
|
||||
agentTarget &&
|
||||
shouldPersistSessionAgentConfiguration(latest.agentConfiguration, agentTarget)
|
||||
) {
|
||||
latest = await sessionPersistenceCoordinator.saveSession({
|
||||
...latest,
|
||||
agentConfiguration: toSessionAgentConfiguration(agentTarget)
|
||||
})
|
||||
}
|
||||
if (!runtime.hasLiveSession(latest.projectId, latest.id) || agentTarget) {
|
||||
await runtime.resumeSession({
|
||||
sessionId: latest.id,
|
||||
cwd: latest.cwd,
|
||||
projectId: latest.projectId,
|
||||
...(latest.permissionProfile
|
||||
? { permissionProfile: latest.permissionProfile }
|
||||
: {}),
|
||||
memoryEnabled: latest.memoryEnabled !== false,
|
||||
...(latest.agentFrameworkId
|
||||
? { previousFrameworkId: latest.agentFrameworkId }
|
||||
: {}),
|
||||
...(latest.agentBackendId ? { previousBackendId: latest.agentBackendId } : {}),
|
||||
...(latest.specialistId ? { specialistId: latest.specialistId } : {}),
|
||||
...(latest.specialistBindingPending === true
|
||||
? { specialistBindingPending: true }
|
||||
: {}),
|
||||
...(latest.providerSessionId
|
||||
? { providerSessionId: latest.providerSessionId }
|
||||
: {}),
|
||||
...(latest.providerContinuityToken
|
||||
? { providerContinuityToken: latest.providerContinuityToken }
|
||||
: {}),
|
||||
...(agentTarget ? { agentTarget } : {})
|
||||
})
|
||||
}
|
||||
const started = await delivery.startDispatch()
|
||||
if (started !== 'started') {
|
||||
throw new DelegateMessageParkedError(
|
||||
'Parent message dispatch fence was not acquired.'
|
||||
)
|
||||
}
|
||||
},
|
||||
delivery.messageId,
|
||||
delivery.onRootAdmissionQueued,
|
||||
(operation) =>
|
||||
archiveCoordinator.withProjectDeletionAdmission(delivery.session.projectId, operation)
|
||||
)
|
||||
}
|
||||
}
|
||||
|
||||
@@ -117,7 +117,7 @@ const create = async (): Promise<string> =>
|
||||
const refresh = async (id: string): Promise<void> => {
|
||||
await owner.execute({ kind: 'smart-collection', collectionId: id, action: 'refresh', offset: 0 })
|
||||
await vi.waitFor(async () => expect((await owner.view(id)).run?.state).toBe('completed'), {
|
||||
timeout: 5000
|
||||
timeout: 15000
|
||||
})
|
||||
}
|
||||
|
||||
|
||||
Reference in New Issue
Block a user