diff --git a/apps/desktop/e2e/workhub-reconstruction.spec.ts b/apps/desktop/e2e/workhub-reconstruction.spec.ts index 938df54eb7..0b29525c03 100644 --- a/apps/desktop/e2e/workhub-reconstruction.spec.ts +++ b/apps/desktop/e2e/workhub-reconstruction.spec.ts @@ -19,7 +19,7 @@ import { COMPOSER_INPUT, ensureSidebarExpanded, expect, test } from './fixtures'; -test('WorkHub rebuilds Session conversation after navigating away and back', async ({ +test('WorkHub rebuilds delegated execution feedback after navigating away and back', async ({ window: page, }) => { const initialPrompt = '检查支付回调重复投递时的幂等性'; @@ -64,6 +64,10 @@ test('WorkHub rebuilds Session conversation after navigating away and back', asy hasText: routedPrompt, }), ).toBeVisible(); + await expect( + page.locator('.workhub-projected-turn', { hasText: routedPrompt }) + .locator('.workhub-submitted-state'), + ).toHaveText('进行中'); }); test('WorkHub defers destructive correction until linked delegation exists', async ({ diff --git a/apps/desktop/src/main/__tests__/runtime-host-session-execution-ipc-main.test.ts b/apps/desktop/src/main/__tests__/runtime-host-session-execution-ipc-main.test.ts index b95338d162..8bdf462bed 100644 --- a/apps/desktop/src/main/__tests__/runtime-host-session-execution-ipc-main.test.ts +++ b/apps/desktop/src/main/__tests__/runtime-host-session-execution-ipc-main.test.ts @@ -630,6 +630,45 @@ test('returns Host-owned cancellation proof to the renderer', async () => { ); }); +test('returns Host-owned Message execution resolutions to the renderer', async () => { + const ipc = ipcHarness(); + registerExecutionIpc( + { + client: executionClient({ + queryMessageExecutions: async (input) => ({ + resolutions: input.messageIds.map((messageId) => messageId === 'message-cancelled' + ? { messageId, state: 'cancelled' as const } + : { + messageId, + state: 'owned' as const, + turnId: 'successor-turn', + runId: 'successor-run', + }), + }), + }), + }, + ipc, + ); + + assert.deepEqual( + await ipc.invoke('sessions:queryMessageExecutions', 'session-1', [ + 'message-delegated', + 'message-cancelled', + ]), + { + resolutions: [ + { + messageId: 'message-delegated', + state: 'owned', + turnId: 'successor-turn', + runId: 'successor-run', + }, + { messageId: 'message-cancelled', state: 'cancelled' }, + ], + }, + ); +}); + test('submits a slash Skill message and reports the Host Skill outcome', async () => { const submits: unknown[] = []; const ipc = ipcHarness(); @@ -1502,6 +1541,7 @@ function executionClient(overrides: Partial): ExecutionClient { interruptTurn: unavailable, listSessionTurnLandmarks: unavailable, listSessionTurns: unavailable, + queryMessageExecutions: unavailable, queryMessages: unavailable, queryTurnResume: unavailable, readExecutionBoundary: unavailable, diff --git a/apps/desktop/src/main/__tests__/workhub-controller.test.ts b/apps/desktop/src/main/__tests__/workhub-controller.test.ts index ff438b0540..81e7538696 100644 --- a/apps/desktop/src/main/__tests__/workhub-controller.test.ts +++ b/apps/desktop/src/main/__tests__/workhub-controller.test.ts @@ -21,12 +21,11 @@ import assert from 'node:assert/strict'; import { existsSync, readFileSync } from 'node:fs'; import test from 'node:test'; import { - createLegacyWorkHubControllerForTests as createWorkHubController, createWorkHubController as createGatedWorkHubController, - WorkHubSessionSubmitError, WORKHUB_ROUTING_STRATEGY_ID, type WorkHubSessionFacts, type WorkHubSessionPort, + type WorkHubCoordinationTurn, } from '../../renderer/workhub-controller.js'; const appShellUrl = [ @@ -71,25 +70,217 @@ function session( }; } -function port(sessions: WorkHubSessionFacts[]): WorkHubSessionPort { +interface TestSessionPort extends WorkHubSessionPort { + create(input: { name: string }): Promise; + submit( + target: { sessionId: string }, + text: string, + turnId: string, + ): Promise<{ turnId: string; steered?: true }>; +} + +function port(sessions: WorkHubSessionFacts[]): TestSessionPort { let nextTurnId = 0; return { list: async () => sessions, recentTurns: async () => [], + delegationFeedback: async (references) => + references.map(({ delegationId }) => ({ delegationId, state: 'accepted' })), routingEvidence: async () => [], create: async () => { throw new Error('create is not used by this read test'); }, - reserveTurnId: () => `reserved-turn-${++nextTurnId}`, - submit: async () => { - throw new Error('submit is not used by this read test'); - }, - reconcileSubmission: async () => ({ kind: 'unknown' }), - stop: async () => {}, + submit: async (_target, _text, turnId) => ({ + turnId: turnId || `reserved-turn-${++nextTurnId}`, + }), subscribe: () => () => {}, }; } +function createWorkHubController({ sessions }: { sessions: TestSessionPort }) { + let candidateByRef = new Map(); + return createGatedWorkHubController({ + sessions, + coordination: { + open: async () => ({ close: async () => undefined }), + record: async (input) => ({ turnId: input.turnId }), + candidates: async () => { + const candidates = (await sessions.list()) + .filter((entry) => entry.kind === 'ordinary' && !entry.archived) + .map((entry) => ({ + candidateRef: `candidate-${entry.target.sessionId}`, + sessionId: entry.target.sessionId, + sessionName: entry.sessionName, + workspace: { + target: { kind: 'host_path' as const, path: `/workspace/${entry.target.sessionId}` }, + hostCwd: `/workspace/${entry.target.sessionId}`, + }, + state: entry.state, + updatedAt: entry.updatedAt, + })); + const byId = new Map( + (await sessions.list()).map((entry) => [entry.target.sessionId, entry]), + ); + candidateByRef = new Map(candidates.flatMap((candidate) => { + const entry = byId.get(candidate.sessionId); + return entry ? [[candidate.candidateRef, entry] as const] : []; + })); + return { + candidateSetId: `sha256:${'a'.repeat(64)}`, + candidates, + }; + }, + act: async (input) => { + if (input.proposal.disposition === 'answer_here') { + return { + disposition: 'answer_here', + coordinationTurnId: input.actionId, + }; + } + if (input.proposal.disposition === 'clarify') { + return { + disposition: 'clarify', + coordinationTurnId: input.actionId, + }; + } + if (input.proposal.disposition === 'create_new') { + const created = await sessions.create({ name: input.proposal.title }); + const admitted = await sessions.submit(created.target, input.userText, input.actionId); + return { + disposition: 'create_new', + targetSessionId: created.target.sessionId, + targetTurnId: admitted.turnId, + ...(admitted.steered ? { steered: true as const } : {}), + }; + } + const target = candidateByRef.get(input.proposal.candidateRef); + if (!target) throw new Error('unknown test candidate'); + const admitted = await sessions.submit(target.target, input.userText, input.actionId); + return { + disposition: 'delegate_existing', + targetSessionId: target.target.sessionId, + targetTurnId: admitted.turnId, + ...(admitted.steered ? { steered: true as const } : {}), + }; + }, + }, + }); +} + +function coordinationAssignmentTurn(): WorkHubCoordinationTurn { + return { + messageId: 'assignment-1', + turnId: 'action-1', + text: 'Continue payments', + state: 'completed', + assignment: { + delegationId: 'delegation-1', + targetSessionId: 'payment', + targetSessionName: 'Payments', + targetMessageId: 'payment-message', + targetTurnId: 'payment-turn', + feedbackState: 'accepted', + }, + updatedAt: 10, + }; +} + +test('conversation acknowledges a durable assignment before projecting target execution', async () => { + const sessions = port([session('payment')]); + let onSessionChanged: (() => void) | undefined; + let feedbackState: 'completed' | 'waiting_for_user' = 'completed'; + sessions.subscribe = (handler) => { + onSessionChanged = handler; + return () => { + onSessionChanged = undefined; + }; + }; + sessions.delegationFeedback = async (references) => + references.map(({ delegationId }) => ({ delegationId, state: feedbackState })); + const assignment = coordinationAssignmentTurn(); + const snapshots: string[] = []; + const controller = createGatedWorkHubController({ + sessions, + coordination: { + open: async (handler) => { + handler([assignment]); + return { close: async () => undefined }; + }, + record: async (input) => ({ turnId: input.turnId }), + candidates: async () => ({ candidateSetId: `sha256:${'a'.repeat(64)}`, candidates: [] }), + act: async () => ({ disposition: 'answer_here', coordinationTurnId: 'unused' }), + }, + }); + + const handle = await controller.openConversation((turns) => { + snapshots.push(turns[0]?.assignment?.feedbackState ?? 'missing'); + }, () => undefined); + await Promise.resolve(); + + assert.deepEqual(snapshots.slice(0, 2), ['accepted', 'completed']); + + feedbackState = 'waiting_for_user'; + onSessionChanged?.(); + await Promise.resolve(); + await Promise.resolve(); + assert.equal(snapshots.at(-1), 'waiting_for_user'); + + await handle.close(); +}); + +test('conversation feedback never lets an older refresh overwrite newer target state', async () => { + const sessions = port([session('payment')]); + let onSessionChanged: (() => void) | undefined; + sessions.subscribe = (handler) => { + onSessionChanged = handler; + return () => undefined; + }; + type Feedback = Awaited>; + const pending: Array<{ + references: Parameters[0]; + resolve(feedback: Feedback): void; + }> = []; + sessions.delegationFeedback = (references) => + new Promise((resolve) => pending.push({ references, resolve })); + const snapshots: string[] = []; + const controller = createGatedWorkHubController({ + sessions, + coordination: { + open: async (handler) => { + handler([coordinationAssignmentTurn()]); + return { close: async () => undefined }; + }, + record: async (input) => ({ turnId: input.turnId }), + candidates: async () => ({ candidateSetId: `sha256:${'b'.repeat(64)}`, candidates: [] }), + act: async () => ({ disposition: 'answer_here', coordinationTurnId: 'unused' }), + }, + }); + + const handle = await controller.openConversation((turns) => { + snapshots.push(turns[0]?.assignment?.feedbackState ?? 'missing'); + }, () => undefined); + assert.equal(pending.length, 1); + onSessionChanged?.(); + assert.equal(pending.length, 2); + + pending[1]!.resolve(pending[1]!.references.map(({ delegationId }) => ({ + delegationId, + state: 'completed', + }))); + await Promise.resolve(); + await Promise.resolve(); + pending[0]!.resolve(pending[0]!.references.map(({ delegationId }) => ({ + delegationId, + state: 'failed', + }))); + await Promise.resolve(); + await Promise.resolve(); + + assert.equal(snapshots.at(-1), 'completed'); + assert.equal(snapshots.includes('failed'), false); + await handle.close(); +}); + test('read exposes existing ordinary Sessions as factual Work summaries', async () => { const controller = createWorkHubController({ sessions: port([ @@ -993,1466 +1184,145 @@ test('English core evidence requires a distinctive word or multiple whole-word m assert.deepEqual(submitted, ['parser']); }); -test('route correction stops the wrong Session and teaches a similar request', async () => { - const submitted: string[] = []; - const stopped: string[] = []; +test('waiting Session rejects a second root request without calling submit', async () => { + let submitted = false; const sessions = port([ - session('login', { sessionName: '登录稳定性' }), - session('payment', { sessionName: '支付稳定性' }), + session('login', { + sessionName: '排查令牌过期重复登录问题', + state: 'waiting_for_user', + }), ]); - sessions.submit = async (target) => { - submitted.push(target.sessionId); - return { turnId: `turn-${submitted.length}` }; - }; - sessions.stop = async (target) => { - stopped.push(target.sessionId); + sessions.submit = async () => { + submitted = true; + return { turnId: 'unexpected' }; }; const controller = createWorkHubController({ sessions }); - await controller.submit({ - requestId: 'request-focus-payment', - text: '先看支付', - explicitTarget: { sessionId: 'payment' }, - }); - - const wrong = await controller.submit({ - requestId: 'request-alias', - text: '继续白鹭点,列出验收项。', - }); - assert.deepEqual(wrong.kind === 'submitted' ? wrong.target : undefined, { - sessionId: 'payment', - }); - const corrected = await controller.submit({ - requestId: 'request-alias', - text: '继续白鹭点,列出验收项。', - explicitTarget: { sessionId: 'login' }, - correction: { from: { sessionId: 'payment' }, turnId: 'turn-2' }, - }); - assert.equal(corrected.kind, 'submitted'); - assert.equal(corrected.kind === 'submitted' ? corrected.evidence : undefined, 'route_correction'); - assert.deepEqual(corrected.kind === 'submitted' ? corrected.correctedFrom : undefined, { - sessionId: 'payment', + const result = await controller.submit({ + requestId: 'request-waiting', + text: '排查令牌过期重复登录问题:补充一条等待状态下的新请求。', }); - const learned = await controller.submit({ - requestId: 'request-alias-similar', - text: '继续白鹭点,补充失败判定。', - }); - assert.deepEqual(learned.kind === 'submitted' ? learned.target : undefined, { - sessionId: 'login', + assert.deepEqual(result, { + kind: 'waiting', + strategyId: WORKHUB_ROUTING_STRATEGY_ID, + requestId: 'request-waiting', + text: '排查令牌过期重复登录问题:补充一条等待状态下的新请求。', + target: { sessionId: 'login' }, }); - assert.equal(learned.kind === 'submitted' ? learned.evidence : undefined, 'route_correction'); - assert.deepEqual(stopped, ['payment']); - assert.deepEqual(submitted, ['payment', 'payment', 'login', 'login']); + assert.equal(submitted, false); }); -test('route correction never stops a root Turn that WorkHub only steered into', async () => { - const stopped: string[] = []; - let submissionCount = 0; +test('submit returns to the previous focused Session', async () => { + const submitted: string[] = []; const sessions = port([ - session('login', { sessionName: '登录稳定性' }), - session('payment', { sessionName: '支付稳定性', state: 'running' }), + session('login', { sessionName: '登录刷新令牌' }), + session('payment', { sessionName: '支付回调幂等性' }), ]); - sessions.submit = async () => { - submissionCount += 1; - return submissionCount === 1 - ? { turnId: 'turn-existing', steered: true } - : { turnId: 'turn-login' }; - }; - sessions.stop = async (target) => { - stopped.push(target.sessionId); + sessions.submit = async (target) => { + submitted.push(target.sessionId); + return { turnId: `turn-${submitted.length}` }; }; const controller = createWorkHubController({ sessions }); - - const wrong = await controller.submit({ - requestId: 'request-steered', - text: '继续补充支付验收项', + await controller.submit({ + requestId: 'request-login', + text: '先看登录', + explicitTarget: { sessionId: 'login' }, + }); + await controller.submit({ + requestId: 'request-payment', + text: '再看支付', explicitTarget: { sessionId: 'payment' }, }); - assert.equal(wrong.kind === 'submitted' ? wrong.steered : undefined, true); - const correction = { - from: { sessionId: 'payment' }, - turnId: 'turn-existing', - steered: true as const, - }; - const corrected = await controller.submit({ - requestId: 'request-steered', - text: '不是支付,应该补充登录验收项', - explicitTarget: { sessionId: 'login' }, - correction, + const result = await controller.submit({ + requestId: 'request-previous', + text: '回到上一个工作', }); - assert.equal(corrected.kind, 'submitted'); - assert.deepEqual(stopped, []); + assert.equal(result.kind, 'submitted'); + assert.deepEqual(result.kind === 'submitted' ? result.target : undefined, { sessionId: 'login' }); + assert.deepEqual(submitted, ['login', 'payment', 'login']); }); -test('first natural-language correction reroutes and stops the wrong WorkHub-owned Turn', async () => { +test('submit lets strong foreign core evidence override a vague focus word', async () => { const submitted: string[] = []; - const stopped: Array<[string, string]> = []; const sessions = port([ session('login', { sessionName: '登录稳定性', - latestResult: '刷新令牌过期导致重复登录', - updatedAt: 20, + latestResult: '处理刷新令牌过期导致的重复登录', }), session('payment', { sessionName: '支付稳定性', - latestResult: '支付回调重复投递', - updatedAt: 30, + latestResult: '处理支付回调重复投递', }), ]); sessions.submit = async (target) => { submitted.push(target.sessionId); return { turnId: `turn-${submitted.length}` }; }; - sessions.stop = async (target, turnId) => { - stopped.push([target.sessionId, turnId]); - }; const controller = createWorkHubController({ sessions }); - await controller.read(); await controller.submit({ - requestId: 'request-wrong-payment', - text: '继续这个工作,补充验收项', + requestId: 'request-payment-focus', + text: '先看支付', + explicitTarget: { sessionId: 'payment' }, }); - const corrected = await controller.submit({ - requestId: 'request-natural-correction', - text: '不是这个,换成登录那个,补充刷新令牌失败判定', + const result = await controller.submit({ + requestId: 'request-foreign-core', + text: '继续处理刷新令牌过期', }); - assert.equal(corrected.kind, 'submitted'); - assert.deepEqual(corrected.kind === 'submitted' ? corrected.target : undefined, { - sessionId: 'login', - }); - assert.equal( - corrected.kind === 'submitted' ? corrected.evidence : undefined, - 'route_correction', - ); - assert.deepEqual(corrected.kind === 'submitted' ? corrected.correctedFrom : undefined, { - sessionId: 'payment', - }); - assert.deepEqual(stopped, [['payment', 'turn-1']]); + assert.equal(result.kind, 'submitted'); + assert.deepEqual(result.kind === 'submitted' ? result.target : undefined, { sessionId: 'login' }); assert.deepEqual(submitted, ['payment', 'login']); }); -test('content-level replacement instructions stay inside the focused Session', async () => { - const submitted: string[] = []; - const stopped: string[] = []; - const sessions = port([ - session('login', { sessionName: '登录稳定性', updatedAt: 20 }), - session('payment', { sessionName: '支付稳定性', updatedAt: 30 }), - session('database', { - sessionName: '数据库迁移', - latestResult: 'Postgres schema migration', - updatedAt: 10, - }), - ]); - sessions.submit = async (target) => { - submitted.push(target.sessionId); - return { turnId: `turn-${submitted.length}` }; - }; - sessions.stop = async (target) => { - stopped.push(target.sessionId); +test('submit keeps unmatched non-executable conversation in WorkHub', async () => { + let created = false; + const actions: unknown[] = []; + const sessions = port([]); + sessions.create = async () => { + created = true; + return session('unexpected'); }; - const controller = createWorkHubController({ sessions }); - await controller.read(); - await controller.submit({ - requestId: 'request-before-content-change', - text: '继续这个工作', + const controller = createGatedWorkHubController({ + sessions, + coordination: { + open: async () => ({ close: async () => undefined }), + record: async (input) => ({ turnId: input.turnId }), + candidates: async () => ({ + candidateSetId: `sha256:${'a'.repeat(64)}`, + candidates: [], + }), + act: async (input) => { + actions.push(input); + return { + disposition: 'answer_here', + coordinationTurnId: 'coordination-turn', + }; + }, + }, }); const result = await controller.submit({ - requestId: 'request-content-change', - text: '继续这个工作,Redis 配置不对,改成 Postgres', + requestId: 'request-discussion', + text: '你觉得统一入口最重要的价值是什么?', }); - assert.deepEqual(result.kind === 'submitted' ? result.target : undefined, { - sessionId: 'payment', + assert.deepEqual(result, { + kind: 'discussion', + strategyId: WORKHUB_ROUTING_STRATEGY_ID, + requestId: 'request-discussion', + text: '你觉得统一入口最重要的价值是什么?', }); - assert.equal(result.kind === 'submitted' ? result.evidence : undefined, 'recent_focus'); - assert.deepEqual(stopped, []); -}); - -test('steering the same WorkHub-owned root preserves ownership for a later correction', async () => { - const submitted: string[] = []; - const stopped: Array<[string, string]> = []; - const sessions = port([ - session('login', { sessionName: '登录稳定性', updatedAt: 20 }), - session('payment', { sessionName: '支付稳定性', updatedAt: 30 }), - ]); - let paymentSubmissions = 0; - sessions.submit = async (target) => { - submitted.push(target.sessionId); - if (target.sessionId === 'payment') { - paymentSubmissions += 1; - return paymentSubmissions === 1 - ? { turnId: 'turn-payment-root' } - : { turnId: 'turn-payment-steering-command', steered: true }; - } - return { turnId: 'turn-login-' + submitted.length }; - }; - sessions.stop = async (target, turnId) => { - stopped.push([target.sessionId, turnId]); - }; - const controller = createWorkHubController({ sessions }); - await controller.read(); - await controller.submit({ - requestId: 'request-owned-root', - text: '继续这个工作', - }); - await controller.submit({ - requestId: 'request-other-owned-root', - text: '先检查登录稳定性', - explicitTarget: { sessionId: 'login' }, - }); - await controller.submit({ - requestId: 'request-steer-owned-root', - text: '继续这个工作,补充测试点', - explicitTarget: { sessionId: 'payment' }, - }); - - const corrected = await controller.submit({ - requestId: 'request-correct-owned-root', - text: '不是这个工作,换成登录稳定性', - }); - - assert.deepEqual(corrected.kind === 'submitted' ? corrected.target : undefined, { - sessionId: 'login', - }); - assert.deepEqual(stopped, [['payment', 'turn-payment-root']]); -}); - -test('a late root completion cannot overwrite newer ownership after remount', async () => { - const stopped: Array<[string, string]> = []; - const sessions = port([ - session('login', { sessionName: '登录稳定性', updatedAt: 20 }), - session('payment', { sessionName: '支付稳定性', updatedAt: 30 }), - ]); - let signalOlderStarted!: () => void; - const olderStarted = new Promise((resolve) => { - signalOlderStarted = resolve; - }); - let finishOlder!: (value: { turnId: string }) => void; - const olderTurn = new Promise<{ turnId: string }>((resolve) => { - finishOlder = resolve; - }); - let paymentSubmissions = 0; - sessions.submit = async (target) => { - if (target.sessionId === 'payment') { - paymentSubmissions += 1; - if (paymentSubmissions === 1) { - signalOlderStarted(); - return olderTurn; - } - return { turnId: 'turn-payment-new' }; - } - return { turnId: 'turn-login' }; - }; - sessions.stop = async (target, turnId) => { - stopped.push([target.sessionId, turnId]); - }; - const controller = createWorkHubController({ sessions }); - - const olderSubmission = controller.submit({ - requestId: 'request-payment-old', - text: '先继续支付稳定性', - explicitTarget: { sessionId: 'payment' }, - }); - await olderStarted; - controller.resetVisitContext(); - await controller.submit({ - requestId: 'request-payment-new', - text: '重新继续支付稳定性', - explicitTarget: { sessionId: 'payment' }, - }); - finishOlder({ turnId: 'turn-payment-old' }); - await olderSubmission; - - const corrected = await controller.submit({ - requestId: 'request-correct-after-late-root', - text: '不是这个工作,换成登录稳定性', - }); - - assert.deepEqual(corrected.kind === 'submitted' ? corrected.target : undefined, { - sessionId: 'login', - }); - assert.deepEqual(stopped, [['payment', 'turn-payment-new']]); -}); - -test('a correction after remount stops a root whose admission is still pending', async () => { - const stopped: Array<[string, string]> = []; - const sessions = port([ - session('login', { sessionName: '登录稳定性', updatedAt: 20 }), - session('payment', { sessionName: '支付稳定性', updatedAt: 30 }), - ]); - let signalPaymentStarted!: () => void; - const paymentStarted = new Promise((resolve) => { - signalPaymentStarted = resolve; - }); - let finishPayment!: (value: { turnId: string }) => void; - const paymentTurn = new Promise<{ turnId: string }>((resolve) => { - finishPayment = resolve; - }); - let nextReservedTurnId = 0; - sessions.reserveTurnId = () => `turn-reserved-${++nextReservedTurnId}`; - sessions.submit = async (target, _text, turnId) => { - if (target.sessionId === 'payment') { - assert.equal(turnId, 'turn-reserved-1'); - signalPaymentStarted(); - return paymentTurn; - } - return { turnId }; - }; - sessions.stop = async (target, turnId) => { - stopped.push([target.sessionId, turnId]); - }; - const controller = createWorkHubController({ sessions }); - - const pendingSubmission = controller.submit({ - requestId: 'request-payment-pending', - text: '继续支付稳定性', - explicitTarget: { sessionId: 'payment' }, - }); - await paymentStarted; - controller.resetVisitContext(); - await controller.read({ focus: { sessionId: 'payment' } }); - - const corrected = await controller.submit({ - requestId: 'request-correct-pending-root', - text: '不是这个工作,换成登录稳定性', - }); - - assert.deepEqual(corrected.kind === 'submitted' ? corrected.target : undefined, { - sessionId: 'login', - }); - assert.deepEqual(stopped, [['payment', 'turn-reserved-1']]); - finishPayment({ turnId: 'turn-payment-host-rebound' }); - await pendingSubmission; - assert.deepEqual(stopped, [ - ['payment', 'turn-reserved-1'], - ['payment', 'turn-payment-host-rebound'], - ]); -}); - -test('a correction stops an uncertain root under the Turn identity the Host minted', async () => { - const stopped: Array<[string, string]> = []; - const sessions = port([ - session('login', { sessionName: '登录稳定性', updatedAt: 20 }), - session('payment', { sessionName: '支付稳定性', updatedAt: 30 }), - ]); - sessions.reserveTurnId = () => 'reserved-payment'; - // The Host admitted the Message and opened a Turn under its own identity, - // then the answer was lost. Only the transcript can tie the reserved Message - // identity back to that Turn. - sessions.submit = async (target, _text, turnId) => { - if (target.sessionId !== 'payment') return { turnId }; - throw new WorkHubSessionSubmitError('delivery outcome is unknown', 'unknown'); - }; - // The transcript has not caught up at delivery time, so the candidate stays - // uncertain; the correction is the next chance to resolve it. - let transcriptCaughtUp = false; - sessions.reconcileSubmission = async (_target, reservedTurnId) => { - if (reservedTurnId !== 'reserved-payment' || !transcriptCaughtUp) { - transcriptCaughtUp = true; - return { kind: 'unknown' }; - } - return { kind: 'root', turnId: 'turn-payment-host' }; - }; - sessions.stop = async (target, turnId) => { - stopped.push([target.sessionId, turnId]); - }; - const controller = createWorkHubController({ sessions }); - - await assert.rejects(controller.submit({ - requestId: 'request-payment-uncertain', - text: '继续支付稳定性', - explicitTarget: { sessionId: 'payment' }, - })); - - const corrected = await controller.submit({ - requestId: 'request-correct-uncertain-root', - text: '不是这个工作,换成登录稳定性', - }); - - assert.deepEqual(corrected.kind === 'submitted' ? corrected.target : undefined, { - sessionId: 'login', - }); - assert.deepEqual(stopped, [['payment', 'turn-payment-host']]); -}); - -test('a correction retries Stop when the same reserved root is admitted before Stop settles', async () => { - const stopped: Array<[string, string]> = []; - const sessions = port([ - session('login', { sessionName: '登录稳定性', updatedAt: 20 }), - session('payment', { sessionName: '支付稳定性', updatedAt: 30 }), - ]); - let signalPaymentStarted!: () => void; - const paymentStarted = new Promise((resolve) => { - signalPaymentStarted = resolve; - }); - let finishPayment!: (value: { turnId: string }) => void; - const paymentTurn = new Promise<{ turnId: string }>((resolve) => { - finishPayment = resolve; - }); - sessions.reserveTurnId = () => 'turn-reserved-1'; - sessions.submit = async (target, _text, turnId) => { - if (target.sessionId === 'payment') { - assert.equal(turnId, 'turn-reserved-1'); - signalPaymentStarted(); - return paymentTurn; - } - return { turnId }; - }; - let signalFirstStopStarted!: () => void; - const firstStopStarted = new Promise((resolve) => { - signalFirstStopStarted = resolve; - }); - let finishFirstStop!: () => void; - const firstStop = new Promise((resolve) => { - finishFirstStop = resolve; - }); - sessions.stop = async (target, turnId) => { - stopped.push([target.sessionId, turnId]); - if (stopped.length === 1) { - signalFirstStopStarted(); - await firstStop; - } - }; - const controller = createWorkHubController({ sessions }); - - const pendingSubmission = controller.submit({ - requestId: 'request-payment-same-id-pending', - text: '继续支付稳定性', - explicitTarget: { sessionId: 'payment' }, - }); - await paymentStarted; - controller.resetVisitContext(); - await controller.read({ focus: { sessionId: 'payment' } }); - - const correction = controller.submit({ - requestId: 'request-correct-same-id-pending-root', - text: '不是这个工作,换成登录稳定性', - }); - await firstStopStarted; - assert.deepEqual(stopped, [['payment', 'turn-reserved-1']]); - - finishPayment({ turnId: 'turn-reserved-1' }); - finishFirstStop(); - await Promise.all([pendingSubmission, correction]); - - assert.deepEqual(stopped, [ - ['payment', 'turn-reserved-1'], - ['payment', 'turn-reserved-1'], - ]); -}); - -test('a stopped ownership tombstone blocks an older root completion', async () => { - const stopped: Array<[string, string]> = []; - const sessions = port([ - session('login', { sessionName: '登录稳定性', updatedAt: 20 }), - session('payment', { sessionName: '支付稳定性', updatedAt: 30 }), - ]); - let signalStaleStarted!: () => void; - const staleStarted = new Promise((resolve) => { - signalStaleStarted = resolve; - }); - let finishStale!: (value: { turnId: string }) => void; - const staleTurn = new Promise<{ turnId: string }>((resolve) => { - finishStale = resolve; - }); - let paymentSubmissions = 0; - sessions.submit = async (target) => { - if (target.sessionId !== 'payment') return { turnId: 'turn-login' }; - paymentSubmissions += 1; - if (paymentSubmissions === 1) return { turnId: 'turn-payment-root' }; - signalStaleStarted(); - return staleTurn; - }; - sessions.stop = async (target, turnId) => { - stopped.push([target.sessionId, turnId]); - }; - const controller = createWorkHubController({ sessions }); - await controller.submit({ - requestId: 'request-payment-owned', - text: '继续支付稳定性', - explicitTarget: { sessionId: 'payment' }, - }); - const staleSubmission = controller.submit({ - requestId: 'request-payment-stale', - text: '再继续支付稳定性', - explicitTarget: { sessionId: 'payment' }, - }); - await staleStarted; - - await controller.submit({ - requestId: 'request-stop-before-stale-finishes', - text: '不是这个工作,换成登录稳定性', - }); - finishStale({ turnId: 'turn-payment-stale' }); - await staleSubmission; - controller.resetVisitContext(); - await controller.read({ focus: { sessionId: 'payment' } }); - await controller.submit({ - requestId: 'request-correct-after-stale-finishes', - text: '不是这个工作,换成登录稳定性', - }); - - assert.deepEqual(stopped, [ - ['payment', 'turn-payment-root'], - ['payment', 'reserved-turn-2'], - ['payment', 'turn-payment-stale'], - ]); -}); - -test('tombstone retention never evicts live ownership for another Session', async () => { - const stopped: Array<[string, string]> = []; - const fillers = Array.from({ length: 32 }, (_, index) => - session(`filler-${index}`, { sessionName: `填充工作 ${index}` })); - const sessions = port([ - session('long-running', { sessionName: '长期工作', updatedAt: 100 }), - session('sink', { sessionName: '收件箱工作', updatedAt: 90 }), - ...fillers, - ]); - sessions.submit = async (target, _text, turnId) => target.sessionId === 'sink' - ? { turnId, steered: true } - : { turnId }; - sessions.stop = async (target, turnId) => { - stopped.push([target.sessionId, turnId]); - }; - const controller = createWorkHubController({ sessions }); - - const live = await controller.submit({ - requestId: 'request-live-root', - text: '开始长期工作', - explicitTarget: { sessionId: 'long-running' }, - }); - assert.equal(live.kind, 'submitted'); - - for (const [index, filler] of fillers.entries()) { - const owned = await controller.submit({ - requestId: `request-filler-${index}`, - text: `开始填充工作 ${index}`, - explicitTarget: filler.target, - }); - assert.equal(owned.kind, 'submitted'); - if (owned.kind !== 'submitted') continue; - await controller.submit({ - requestId: `request-stop-filler-${index}`, - text: `填充工作 ${index} 路由错了`, - explicitTarget: { sessionId: 'sink' }, - correction: { - from: filler.target, - turnId: owned.turnId, - }, - }); - } - - stopped.length = 0; - controller.resetVisitContext(); - await controller.read({ focus: { sessionId: 'long-running' } }); - await controller.submit({ - requestId: 'request-correct-live-root', - text: '不是这个工作,换成收件箱工作', - }); - - assert.equal(stopped.length, 1); - assert.equal(stopped[0]?.[0], 'long-running'); - assert.equal(stopped[0]?.[1], live.kind === 'submitted' ? live.turnId : undefined); -}); - -test('correction barrier rejects a new root while Stop is pending', async () => { - const stopped: Array<[string, string]> = []; - const sessions = port([ - session('login', { sessionName: '登录稳定性', updatedAt: 20 }), - session('payment', { sessionName: '支付稳定性', updatedAt: 30 }), - ]); - let signalStopStarted!: () => void; - const stopStarted = new Promise((resolve) => { - signalStopStarted = resolve; - }); - let finishFirstStop!: () => void; - const firstStop = new Promise((resolve) => { - finishFirstStop = resolve; - }); - let stopCalls = 0; - sessions.submit = async (_target, _text, turnId) => ({ turnId }); - sessions.stop = async (target, turnId) => { - stopped.push([target.sessionId, turnId]); - stopCalls += 1; - if (stopCalls === 1) { - signalStopStarted(); - await firstStop; - } - }; - const controller = createWorkHubController({ sessions }); - const original = await controller.submit({ - requestId: 'request-original-payment', - text: '继续支付稳定性', - explicitTarget: { sessionId: 'payment' }, - }); - assert.equal(original.kind, 'submitted'); - - const correction = controller.submit({ - requestId: 'request-correct-original-payment', - text: '不是这个工作,换成登录稳定性', - }); - await stopStarted; - await assert.rejects(controller.submit({ - requestId: 'request-overlapping-payment', - text: '重新处理支付稳定性', - explicitTarget: { sessionId: 'payment' }, - }), /still reconciling/u); - finishFirstStop(); - await correction; - - const newer = await controller.submit({ - requestId: 'request-newer-payment', - text: '重新处理支付稳定性', - explicitTarget: { sessionId: 'payment' }, - }); - assert.equal(newer.kind, 'submitted'); - - controller.resetVisitContext(); - await controller.read({ focus: { sessionId: 'payment' } }); - await controller.submit({ - requestId: 'request-correct-newer-payment', - text: '不是这个工作,换成登录稳定性', - }); - - assert.deepEqual(stopped, [ - ['payment', original.kind === 'submitted' ? original.turnId : ''], - ['payment', newer.kind === 'submitted' ? newer.turnId : ''], - ]); -}); - -test('a partially failed correction records only successful Stops and keeps failures reachable', async () => { - const stopAttempts: Array<[string, string]> = []; - const sessions = port([ - session('login', { sessionName: '登录稳定性', updatedAt: 20 }), - session('payment', { sessionName: '支付稳定性', updatedAt: 30 }), - ]); - let finishPending!: (value: { turnId: string }) => void; - const pendingTurn = new Promise<{ turnId: string }>((resolve) => { - finishPending = resolve; - }); - let paymentSubmissions = 0; - sessions.reserveTurnId = (() => { - let next = 0; - return () => `reserved-${++next}`; - })(); - sessions.submit = async (target, _text, turnId) => { - if (target.sessionId !== 'payment') return { turnId }; - paymentSubmissions += 1; - return paymentSubmissions === 1 ? { turnId: 'payment-root' } : pendingTurn; - }; - let failPaymentRoot = true; - sessions.stop = async (target, turnId) => { - stopAttempts.push([target.sessionId, turnId]); - if (turnId === 'payment-root' && failPaymentRoot) { - failPaymentRoot = false; - throw new Error('Host rejected the first Stop'); - } - }; - const controller = createWorkHubController({ sessions }); - const confirmed = await controller.submit({ - requestId: 'confirmed-payment-root', - text: '继续支付稳定性', - explicitTarget: { sessionId: 'payment' }, - }); - assert.equal(confirmed.kind, 'submitted'); - const pending = controller.submit({ - requestId: 'pending-payment-root', - text: '再次处理支付稳定性', - explicitTarget: { sessionId: 'payment' }, - }); - await Promise.resolve(); - await Promise.resolve(); - - await assert.rejects(controller.submit({ - requestId: 'partially-failed-correction', - text: '不是支付,改成登录', - explicitTarget: { sessionId: 'login' }, - correction: { - from: { sessionId: 'payment' }, - turnId: 'payment-root', - }, - }), /Host rejected the first Stop/u); - - assert.equal(stopAttempts.some(([, turnId]) => turnId === 'reserved-2'), true); - finishPending({ turnId: 'payment-host-rebound' }); - await pending; - assert.equal(stopAttempts.some(([, turnId]) => turnId === 'payment-host-rebound'), true); - - await controller.submit({ - requestId: 'retry-failed-correction', - text: '不是支付,改成登录', - explicitTarget: { sessionId: 'login' }, - correction: { - from: { sessionId: 'payment' }, - turnId: 'payment-root', - }, - }); - assert.equal( - stopAttempts.filter(([, turnId]) => turnId === 'payment-root').length, - 2, - ); -}); - -test('the 33rd unresolved root is back-pressured before Host admission', async () => { - const facts = Array.from({ length: 33 }, (_, index) => - session(`work-${index}`, { runningTurnIds: [] })); - const sessions = port(facts); - const finishes = new Map void>(); - let admitted = 0; - let signalThirtyTwo!: () => void; - const thirtyTwoAdmitted = new Promise((resolve) => { - signalThirtyTwo = resolve; - }); - sessions.submit = async (target, _text, turnId) => { - admitted += 1; - if (admitted === 32) signalThirtyTwo(); - return await new Promise<{ turnId: string }>((resolve) => { - finishes.set(target.sessionId, resolve); - }); - }; - const controller = createWorkHubController({ sessions }); - const inFlight = facts.slice(0, 32).map((fact, index) => controller.submit({ - requestId: `root-${index}`, - text: `开始工作 ${index}`, - explicitTarget: fact.target, - })); - await thirtyTwoAdmitted; - - await assert.rejects(controller.submit({ - requestId: 'root-33', - text: '开始工作 33', - explicitTarget: facts[32]!.target, - }), /too many unresolved root submissions/u); - assert.equal(admitted, 32); - - finishes.get('work-0')?.({ turnId: 'settled-work-0' }); - await inFlight[0]; - const thirtyThird = controller.submit({ - requestId: 'root-33-after-capacity', - text: '开始工作 33', - explicitTarget: facts[32]!.target, - }); - await Promise.resolve(); - await Promise.resolve(); - assert.equal(admitted, 33); - finishes.get('work-32')?.({ turnId: 'settled-work-32' }); - await thirtyThird; -}); - -test('an old correction barrier survives more than 32 newer corrections', async () => { - const stopped: Array<[string, string]> = []; - const fillers = Array.from({ length: 32 }, (_, index) => - session(`barrier-filler-${index}`, { sessionName: `屏障填充 ${index}` })); - const sessions = port([ - session('old-pending', { sessionName: '旧的待定工作', updatedAt: 100 }), - session('sink', { sessionName: '安全收件箱', updatedAt: 90 }), - ...fillers, - ]); - let finishOld!: (value: { turnId: string }) => void; - const oldTurn = new Promise<{ turnId: string }>((resolve) => { - finishOld = resolve; - }); - sessions.submit = async (target, _text, turnId) => { - if (target.sessionId === 'old-pending') return oldTurn; - if (target.sessionId === 'sink') return { turnId, steered: true }; - return { turnId }; - }; - sessions.stop = async (target, turnId) => { - stopped.push([target.sessionId, turnId]); - }; - const controller = createWorkHubController({ sessions }); - const oldSubmission = controller.submit({ - requestId: 'old-pending-root', - text: '开始旧的待定工作', - explicitTarget: { sessionId: 'old-pending' }, - }); - await Promise.resolve(); - await Promise.resolve(); - await controller.submit({ - requestId: 'correct-old-pending', - text: '旧工作路由错了', - explicitTarget: { sessionId: 'sink' }, - correction: { - from: { sessionId: 'old-pending' }, - turnId: 'reserved-turn-1', - }, - }); - - for (const [index, filler] of fillers.entries()) { - const owned = await controller.submit({ - requestId: `newer-barrier-root-${index}`, - text: `开始屏障填充 ${index}`, - explicitTarget: filler.target, - }); - assert.equal(owned.kind, 'submitted'); - if (owned.kind !== 'submitted') continue; - await controller.submit({ - requestId: `newer-barrier-correction-${index}`, - text: `屏障填充 ${index} 路由错了`, - explicitTarget: { sessionId: 'sink' }, - correction: { from: filler.target, turnId: owned.turnId }, - }); - } - - finishOld({ turnId: 'old-host-rebound' }); - await oldSubmission; - assert.equal( - stopped.some(([sessionId, turnId]) => - sessionId === 'old-pending' && turnId === 'old-host-rebound'), - true, - ); -}); - -test('a definite Host rejection releases only its own pending admission', async () => { - const stopped: Array<[string, string]> = []; - const sessions = port([ - session('login', { sessionName: '登录稳定性' }), - session('payment', { sessionName: '支付稳定性' }), - ]); - let shouldReject = true; - sessions.submit = async (_target, _text, turnId) => { - if (shouldReject) { - shouldReject = false; - throw new WorkHubSessionSubmitError('not admitted', 'rejected'); - } - return { turnId }; - }; - sessions.stop = async (target, turnId) => { - stopped.push([target.sessionId, turnId]); - }; - const controller = createWorkHubController({ sessions }); - - await assert.rejects(controller.submit({ - requestId: 'definitely-rejected', - text: '继续支付稳定性', - explicitTarget: { sessionId: 'payment' }, - }), /not admitted/u); - const admitted = await controller.submit({ - requestId: 'later-admitted', - text: '继续支付稳定性', - explicitTarget: { sessionId: 'payment' }, - }); - assert.equal(admitted.kind, 'submitted'); - await controller.submit({ - requestId: 'correct-later-admitted', - text: '不是支付,改成登录', - explicitTarget: { sessionId: 'login' }, - correction: { from: { sessionId: 'payment' } }, - }); - - assert.deepEqual(stopped, [[ - 'payment', - admitted.kind === 'submitted' ? admitted.turnId : '', - ]]); -}); - -test('a lost delivery reply is reconciled to its authoritative root ownership', async () => { - const stopped: Array<[string, string]> = []; - const sessions = port([ - session('login', { sessionName: '登录稳定性' }), - session('payment', { sessionName: '支付稳定性' }), - ]); - sessions.submit = async () => { - throw new WorkHubSessionSubmitError('reply lost', 'unknown'); - }; - sessions.reconcileSubmission = async () => ({ - kind: 'root', - turnId: 'authoritative-payment-root', - }); - sessions.stop = async (target, turnId) => { - stopped.push([target.sessionId, turnId]); - }; - const controller = createWorkHubController({ sessions }); - - await assert.rejects(controller.submit({ - requestId: 'reply-lost', - text: '继续支付稳定性', - explicitTarget: { sessionId: 'payment' }, - }), /reply lost/u); - sessions.submit = async (_target, _text, turnId) => ({ turnId }); - await controller.submit({ - requestId: 'correct-reconciled-root', - text: '不是支付,改成登录', - explicitTarget: { sessionId: 'login' }, - correction: { from: { sessionId: 'payment' } }, - }); - - assert.deepEqual(stopped, [['payment', 'authoritative-payment-root']]); -}); - -test('an unknown delivery remains pending until a later authoritative reconciliation', async () => { - const stopped: Array<[string, string]> = []; - const sessions = port([ - session('login', { sessionName: '登录稳定性' }), - session('payment', { sessionName: '支付稳定性' }), - ]); - sessions.submit = async () => { - throw new WorkHubSessionSubmitError('reply lost', 'unknown'); - }; - let reconciliation: Awaited> = { - kind: 'unknown', - }; - sessions.reconcileSubmission = async () => reconciliation; - sessions.stop = async (target, turnId) => { - stopped.push([target.sessionId, turnId]); - }; - const controller = createWorkHubController({ sessions }); - - await assert.rejects(controller.submit({ - requestId: 'unknown-delivery', - text: '继续支付稳定性', - explicitTarget: { sessionId: 'payment' }, - }), /reply lost/u); - reconciliation = { kind: 'root', turnId: 'later-authoritative-root' }; - await controller.read(); - sessions.submit = async (_target, _text, turnId) => ({ turnId }); - await controller.submit({ - requestId: 'correct-later-reconciled-root', - text: '不是支付,改成登录', - explicitTarget: { sessionId: 'login' }, - correction: { from: { sessionId: 'payment' } }, - }); - - assert.deepEqual(stopped, [['payment', 'later-authoritative-root']]); -}); - -test('a partial multi-Host catalog never erases confirmed ownership', async () => { - const stopped: Array<[string, string]> = []; - const allSessions = [ - session('login', { sessionName: '登录稳定性' }), - session('remote-payment', { sessionName: '远端支付稳定性' }), - ]; - const sessions = port(allSessions); - let visible = allSessions; - let paymentCatalogComplete = true; - sessions.listCatalog = async () => ({ - sessions: visible, - isCompleteFor: (target) => - target.sessionId !== 'remote-payment' || paymentCatalogComplete, - }); - sessions.submit = async (_target, _text, turnId) => ({ turnId }); - sessions.stop = async (target, turnId) => { - stopped.push([target.sessionId, turnId]); - }; - const controller = createWorkHubController({ sessions }); - const owned = await controller.submit({ - requestId: 'remote-root', - text: '继续远端支付', - explicitTarget: { sessionId: 'remote-payment' }, - }); - assert.equal(owned.kind, 'submitted'); - - visible = [allSessions[0]!]; - paymentCatalogComplete = false; - await controller.submit({ - requestId: 'correct-after-partial-catalog', - text: '不是支付,改成登录', - explicitTarget: { sessionId: 'login' }, - correction: { from: { sessionId: 'remote-payment' } }, - }); - - assert.deepEqual(stopped, [[ - 'remote-payment', - owned.kind === 'submitted' ? owned.turnId : '', - ]]); -}); - -test('a stale complete catalog never erases newer confirmed ownership', async () => { - const stopped: Array<[string, string]> = []; - const allSessions = [ - session('login', { sessionName: '登录稳定性' }), - session('payment', { sessionName: '支付稳定性' }), - ]; - const sessions = port(allSessions); - let signalStaleReadStarted!: () => void; - const staleReadStarted = new Promise((resolve) => { - signalStaleReadStarted = resolve; - }); - let finishStaleRead!: (value: { - sessions: WorkHubSessionFacts[]; - isCompleteFor(target: { sessionId: string }): boolean; - }) => void; - const staleCatalog = new Promise<{ - sessions: WorkHubSessionFacts[]; - isCompleteFor(target: { sessionId: string }): boolean; - }>((resolve) => { - finishStaleRead = resolve; - }); - let catalogReads = 0; - sessions.listCatalog = async () => { - catalogReads += 1; - if (catalogReads === 1) { - signalStaleReadStarted(); - return staleCatalog; - } - return { - sessions: allSessions, - isCompleteFor: () => true, - }; - }; - sessions.submit = async (_target, _text, turnId) => ({ turnId }); - sessions.stop = async (target, turnId) => { - stopped.push([target.sessionId, turnId]); - }; - const controller = createWorkHubController({ sessions }); - - const staleRead = controller.read(); - await staleReadStarted; - const owned = await controller.submit({ - requestId: 'payment-root-after-stale-read-started', - text: '继续支付稳定性', - explicitTarget: { sessionId: 'payment' }, - }); - assert.equal(owned.kind, 'submitted'); - - finishStaleRead({ - sessions: [allSessions[0]!], - isCompleteFor: () => true, - }); - await staleRead; - await controller.submit({ - requestId: 'correct-after-stale-complete-catalog', - text: '不是支付,改成登录', - explicitTarget: { sessionId: 'login' }, - correction: { from: { sessionId: 'payment' } }, - }); - - assert.deepEqual(stopped, [[ - 'payment', - owned.kind === 'submitted' ? owned.turnId : '', - ]]); -}); - -test('authoritative Session removal releases uncertain admissions before global backpressure', async () => { - const stale = Array.from({ length: 32 }, (_, index) => - session(`removed-${index}`, { sessionName: `已删除工作 ${index}` })); - const fresh = session('fresh', { sessionName: '新工作' }); - const sessions = port(stale); - let visible = stale; - let removedCatalogComplete = false; - sessions.listCatalog = async () => ({ - sessions: visible, - isCompleteFor: (target) => - target.sessionId.startsWith('removed-') && removedCatalogComplete, - }); - sessions.submit = async () => { - throw new WorkHubSessionSubmitError('reply lost', 'unknown'); - }; - sessions.reconcileSubmission = async () => ({ kind: 'unknown' }); - const controller = createWorkHubController({ sessions }); - - for (const [index, fact] of stale.entries()) { - await assert.rejects(controller.submit({ - requestId: `uncertain-${index}`, - text: `开始已删除工作 ${index}`, - explicitTarget: fact.target, - }), /reply lost/u); - } - - visible = [fresh]; - removedCatalogComplete = true; - sessions.submit = async (_target, _text, turnId) => ({ turnId }); - const admitted = await controller.submit({ - requestId: 'after-authoritative-removal', - text: '开始新工作', - explicitTarget: fresh.target, - }); - - assert.equal(admitted.kind, 'submitted'); - assert.deepEqual(admitted.kind === 'submitted' ? admitted.target : undefined, fresh.target); -}); - -test('a lost reply reconciled as steering never claims the pre-existing root', async () => { - const stopped: Array<[string, string]> = []; - const sessions = port([ - session('login', { sessionName: '登录稳定性' }), - session('payment', { sessionName: '支付稳定性' }), - ]); - sessions.submit = async () => { - throw new WorkHubSessionSubmitError('steering reply lost', 'unknown'); - }; - sessions.reconcileSubmission = async () => ({ kind: 'steered' }); - sessions.stop = async (target, turnId) => { - stopped.push([target.sessionId, turnId]); - }; - const controller = createWorkHubController({ sessions }); - - await assert.rejects(controller.submit({ - requestId: 'steering-reply-lost', - text: '补充支付测试', - explicitTarget: { sessionId: 'payment' }, - }), /steering reply lost/u); - sessions.submit = async (_target, _text, turnId) => ({ turnId }); - await controller.submit({ - requestId: 'correct-after-steering-reply-loss', - text: '不是支付,改成登录', - explicitTarget: { sessionId: 'login' }, - correction: { from: { sessionId: 'payment' } }, - }); - - assert.deepEqual(stopped, []); -}); - -test('WorkHub-owned root remains stoppable after navigating away and back', async () => { - const stopped: Array<[string, string]> = []; - const sessions = port([ - session('login', { sessionName: '登录稳定性', updatedAt: 20 }), - session('payment', { sessionName: '支付稳定性', updatedAt: 30 }), - ]); - sessions.submit = async (target) => ({ turnId: 'turn-' + target.sessionId }); - sessions.stop = async (target, turnId) => { - stopped.push([target.sessionId, turnId]); - }; - const controller = createWorkHubController({ sessions }); - await controller.read(); - await controller.submit({ - requestId: 'request-owned-before-navigation', - text: '继续这个工作', - }); - controller.resetVisitContext(); - await controller.read({ focus: { sessionId: 'payment' } }); - - const corrected = await controller.submit({ - requestId: 'request-correction-after-return', - text: '不是这个工作,换成登录稳定性', - }); - - assert.deepEqual(corrected.kind === 'submitted' ? corrected.target : undefined, { - sessionId: 'login', - }); - assert.deepEqual(stopped, [['payment', 'turn-payment']]); -}); - -test('natural-language correction never stops a pre-existing focused Session', async () => { - const stopped: string[] = []; - const sessions = port([ - session('login', { sessionName: '登录稳定性', updatedAt: 20 }), - session('payment', { sessionName: '支付稳定性', state: 'running', updatedAt: 30 }), - ]); - sessions.submit = async (target) => ({ turnId: `turn-${target.sessionId}` }); - sessions.stop = async (target) => { - stopped.push(target.sessionId); - }; - const controller = createWorkHubController({ sessions }); - await controller.read(); - - const corrected = await controller.submit({ - requestId: 'request-safe-natural-correction', - text: '不是这个,用登录那个', - }); - - assert.deepEqual(corrected.kind === 'submitted' ? corrected.target : undefined, { - sessionId: 'login', - }); - assert.deepEqual(corrected.kind === 'submitted' ? corrected.correctedFrom : undefined, { - sessionId: 'payment', - }); - assert.deepEqual(stopped, []); -}); - -test('English natural-language correction names the replacement Session', async () => { - const submitted: string[] = []; - const sessions = port([ - session('login', { sessionName: 'Login Reliability', updatedAt: 20 }), - session('payment', { sessionName: 'Payment Webhooks', updatedAt: 30 }), - ]); - sessions.submit = async (target) => { - submitted.push(target.sessionId); - return { turnId: `turn-${submitted.length}` }; - }; - const controller = createWorkHubController({ sessions }); - await controller.read(); - - const corrected = await controller.submit({ - requestId: 'request-english-natural-correction', - text: 'Not that work; switch to Login Reliability and add the retry checks', - }); - - assert.deepEqual(corrected.kind === 'submitted' ? corrected.target : undefined, { - sessionId: 'login', - }); - assert.equal( - corrected.kind === 'submitted' ? corrected.evidence : undefined, - 'route_correction', - ); -}); - -test('ambiguous natural-language correction preserves correction context through clarification', async () => { - const submitted: string[] = []; - const stopped: Array<[string, string]> = []; - const sessions = port([ - session('login-api', { sessionName: '登录 API 稳定性', updatedAt: 20 }), - session('login-ui', { sessionName: '登录 UI 稳定性', updatedAt: 10 }), - session('payment', { sessionName: '支付回调幂等性', updatedAt: 30 }), - ]); - sessions.submit = async (target) => { - submitted.push(target.sessionId); - return { turnId: `turn-${submitted.length}` }; - }; - sessions.stop = async (target, turnId) => { - stopped.push([target.sessionId, turnId]); - }; - const controller = createWorkHubController({ sessions }); - await controller.read(); - await controller.submit({ - requestId: 'request-payment-before-clarification', - text: '继续这个工作', - }); - - const clarification = await controller.submit({ - requestId: 'request-natural-clarification', - text: '不是这个,换成登录那个', - }); - assert.equal(clarification.kind, 'clarification'); - if (clarification.kind !== 'clarification') return; - assert.deepEqual( - clarification.options.map((option) => option.target.sessionId), - ['login-api', 'login-ui'], - ); - assert.deepEqual(clarification.correction, { - from: { sessionId: 'payment' }, - turnId: 'turn-1', - }); - - const corrected = await controller.submit({ - requestId: clarification.requestId, - text: clarification.text, - explicitTarget: { sessionId: 'login-api' }, - correction: clarification.correction, - }); - assert.equal(corrected.kind, 'submitted'); - assert.deepEqual(stopped, [['payment', 'turn-1']]); - assert.deepEqual(submitted, ['payment', 'login-api']); -}); - -test('latest route correction wins for the same expression family', async () => { - const sessions = port([ - session('login', { sessionName: '登录稳定性' }), - session('payment', { sessionName: '支付稳定性' }), - ]); - sessions.submit = async (_target) => ({ turnId: 'turn' }); - const controller = createWorkHubController({ sessions }); - - await controller.submit({ - requestId: 'correction-login', - text: '继续白鹭点,列出验收项。', - explicitTarget: { sessionId: 'login' }, - correction: { from: { sessionId: 'payment' }, turnId: 'turn' }, - }); - await controller.submit({ - requestId: 'correction-payment', - text: '继续白鹭点,列出异常项。', - explicitTarget: { sessionId: 'payment' }, - correction: { from: { sessionId: 'login' }, turnId: 'turn' }, - }); - - const result = await controller.submit({ - requestId: 'correction-latest', - text: '继续白鹭点,补充回滚条件。', - }); - - assert.deepEqual(result.kind === 'submitted' ? result.target : undefined, { - sessionId: 'payment', - }); - assert.equal(result.kind === 'submitted' ? result.evidence : undefined, 'route_correction'); -}); - -test('user correction order wins when overlapping submissions finish out of order', async () => { - const sessions = port([ - session('login', { sessionName: '登录稳定性' }), - session('payment', { sessionName: '支付稳定性' }), - ]); - let signalOlderStarted!: () => void; - const olderStarted = new Promise((resolve) => { - signalOlderStarted = resolve; - }); - let finishOlder!: (value: { turnId: string }) => void; - const olderTurn = new Promise<{ turnId: string }>((resolve) => { - finishOlder = resolve; - }); - sessions.submit = async (target) => { - if (target.sessionId === 'login') { - signalOlderStarted(); - return olderTurn; - } - return { turnId: 'turn-payment' }; - }; - const controller = createWorkHubController({ sessions }); - - const olderCorrection = controller.submit({ - requestId: 'correction-older-login', - text: '继续白鹭点,列出验收项。', - explicitTarget: { sessionId: 'login' }, - correction: { from: { sessionId: 'payment' } }, - }); - await olderStarted; - controller.resetVisitContext(); - await controller.submit({ - requestId: 'correction-newer-payment', - text: '继续白鹭点,列出异常项。', - explicitTarget: { sessionId: 'payment' }, - correction: { from: { sessionId: 'login' } }, - }); - finishOlder({ turnId: 'turn-login' }); - await olderCorrection; - - const result = await controller.submit({ - requestId: 'correction-after-overlap', - text: '继续白鹭点,补充回滚条件。', - }); - - assert.deepEqual(result.kind === 'submitted' ? result.target : undefined, { - sessionId: 'payment', - }); - assert.equal(result.kind === 'submitted' ? result.evidence : undefined, 'route_correction'); -}); - -test('waiting Session rejects a second root request without calling submit', async () => { - let submitted = false; - const sessions = port([ - session('login', { - sessionName: '排查令牌过期重复登录问题', - state: 'waiting_for_user', - }), - ]); - sessions.submit = async () => { - submitted = true; - return { turnId: 'unexpected' }; - }; - const controller = createWorkHubController({ sessions }); - - const result = await controller.submit({ - requestId: 'request-waiting', - text: '排查令牌过期重复登录问题:补充一条等待状态下的新请求。', - }); - - assert.deepEqual(result, { - kind: 'waiting', - strategyId: WORKHUB_ROUTING_STRATEGY_ID, - requestId: 'request-waiting', - text: '排查令牌过期重复登录问题:补充一条等待状态下的新请求。', - target: { sessionId: 'login' }, - }); - assert.equal(submitted, false); -}); - -test('submit returns to the previous focused Session', async () => { - const submitted: string[] = []; - const sessions = port([ - session('login', { sessionName: '登录刷新令牌' }), - session('payment', { sessionName: '支付回调幂等性' }), - ]); - sessions.submit = async (target) => { - submitted.push(target.sessionId); - return { turnId: `turn-${submitted.length}` }; - }; - const controller = createWorkHubController({ sessions }); - await controller.submit({ - requestId: 'request-login', - text: '先看登录', - explicitTarget: { sessionId: 'login' }, - }); - await controller.submit({ - requestId: 'request-payment', - text: '再看支付', - explicitTarget: { sessionId: 'payment' }, - }); - - const result = await controller.submit({ - requestId: 'request-previous', - text: '回到上一个工作', - }); - - assert.equal(result.kind, 'submitted'); - assert.deepEqual(result.kind === 'submitted' ? result.target : undefined, { sessionId: 'login' }); - assert.deepEqual(submitted, ['login', 'payment', 'login']); -}); - -test('submit lets strong foreign core evidence override a vague focus word', async () => { - const submitted: string[] = []; - const sessions = port([ - session('login', { - sessionName: '登录稳定性', - latestResult: '处理刷新令牌过期导致的重复登录', - }), - session('payment', { - sessionName: '支付稳定性', - latestResult: '处理支付回调重复投递', - }), - ]); - sessions.submit = async (target) => { - submitted.push(target.sessionId); - return { turnId: `turn-${submitted.length}` }; - }; - const controller = createWorkHubController({ sessions }); - await controller.submit({ - requestId: 'request-payment-focus', - text: '先看支付', - explicitTarget: { sessionId: 'payment' }, - }); - - const result = await controller.submit({ - requestId: 'request-foreign-core', - text: '继续处理刷新令牌过期', - }); - - assert.equal(result.kind, 'submitted'); - assert.deepEqual(result.kind === 'submitted' ? result.target : undefined, { sessionId: 'login' }); - assert.deepEqual(submitted, ['payment', 'login']); -}); - -test('submit keeps unmatched non-executable conversation in WorkHub', async () => { - let created = false; - const actions: unknown[] = []; - const sessions = port([]); - sessions.create = async () => { - created = true; - return session('unexpected'); - }; - const controller = createGatedWorkHubController({ - sessions, - coordination: { - open: async () => ({ close: async () => undefined }), - answer: async (input) => ({ turnId: input.turnId }), - record: async (input) => ({ turnId: input.turnId }), - candidates: async () => ({ - candidateSetId: `sha256:${'a'.repeat(64)}`, - candidates: [], - }), - act: async (input) => { - actions.push(input); - return { - disposition: 'answer_here', - coordinationTurnId: 'coordination-turn', - }; - }, - }, - }); - - const result = await controller.submit({ - requestId: 'request-discussion', - text: '你觉得统一入口最重要的价值是什么?', - }); - - assert.deepEqual(result, { - kind: 'discussion', - strategyId: WORKHUB_ROUTING_STRATEGY_ID, - requestId: 'request-discussion', - text: '你觉得统一入口最重要的价值是什么?', - }); - assert.equal(created, false); - assert.deepEqual(actions, [ - { - actionId: 'request-discussion', - userText: '你觉得统一入口最重要的价值是什么?', - proposal: { disposition: 'answer_here' }, - }, + assert.equal(created, false); + assert.deepEqual(actions, [ + { + actionId: 'request-discussion', + userText: '你觉得统一入口最重要的价值是什么?', + proposal: { disposition: 'answer_here' }, + }, ]); }); @@ -2462,14 +1332,10 @@ test('production submission delegates only through the Runtime-owned candidate r sessions.submit = async () => { throw new Error('renderer direct submit must not be used'); }; - sessions.stop = async () => { - throw new Error('renderer direct stop must not be used'); - }; const controller = createGatedWorkHubController({ sessions, coordination: { open: async () => ({ close: async () => undefined }), - answer: async (input) => ({ turnId: input.turnId }), record: async (input) => ({ turnId: input.turnId }), candidates: async () => ({ candidateSetId: `sha256:${'b'.repeat(64)}`, @@ -2524,7 +1390,6 @@ test('production retry reaches durable Action Gate replay while target is waitin sessions, coordination: { open: async () => ({ close: async () => undefined }), - answer: async (input) => ({ turnId: input.turnId }), record: async (input) => ({ turnId: input.turnId }), candidates: async () => ({ candidateSetId: `sha256:${'c'.repeat(64)}`, @@ -2570,7 +1435,6 @@ test('production defers destructive correction until persistent delegation exist sessions, coordination: { open: async () => ({ close: async () => undefined }), - answer: async (input) => ({ turnId: input.turnId }), record: async (input) => ({ turnId: input.turnId }), candidates: async () => ({ candidateSetId: `sha256:${'d'.repeat(64)}`, @@ -2676,7 +1540,6 @@ test('production natural-language correction fails closed before a second delega sessions, coordination: { open: async () => ({ close: async () => undefined }), - answer: async (input) => ({ turnId: input.turnId }), record: async (input) => ({ turnId: input.turnId }), candidates: async () => ({ candidateSetId, candidates }), act: async (input) => { @@ -2732,7 +1595,6 @@ test('production correction-shaped creation stays create_new without an existing sessions: port([]), coordination: { open: async () => ({ close: async () => undefined }), - answer: async (input) => ({ turnId: input.turnId }), record: async (input) => ({ turnId: input.turnId }), candidates: async () => ({ candidateSetId: `sha256:${'c'.repeat(64)}`, @@ -2764,7 +1626,6 @@ test('production clarification is persisted through the typed Action Gate dispos sessions: port([]), coordination: { open: async () => ({ close: async () => undefined }), - answer: async (input) => ({ turnId: input.turnId }), record: async () => { throw new Error('legacy summary recording must not persist clarification'); }, @@ -2808,7 +1669,6 @@ test('production creation leaves Session identity and workspace authority to mai sessions, coordination: { open: async () => ({ close: async () => undefined }), - answer: async (input) => ({ turnId: input.turnId }), record: async (input) => ({ turnId: input.turnId }), candidates: async () => ({ candidateSetId: `sha256:${'c'.repeat(64)}`, diff --git a/apps/desktop/src/main/__tests__/workhub-coordination-host-scope.test.ts b/apps/desktop/src/main/__tests__/workhub-coordination-host-scope.test.ts index 0a4e014fe2..5feca3844c 100644 --- a/apps/desktop/src/main/__tests__/workhub-coordination-host-scope.test.ts +++ b/apps/desktop/src/main/__tests__/workhub-coordination-host-scope.test.ts @@ -19,10 +19,7 @@ import assert from 'node:assert/strict'; import test from 'node:test'; -import { - desktopSessionKey, - parseDesktopSessionKey, -} from '../../shared/runtime-host-identity.js'; +import { desktopSessionKey } from '../../shared/runtime-host-identity.js'; import { scopeWorkHubSessionsToCoordinationHost } from '../../renderer/workhub-coordination-host-scope.js'; import { startWorkHubCoordinationLifecycle, @@ -30,7 +27,7 @@ import { } from '../../renderer/workhub-coordination-lifecycle.js'; import type { WorkHubDesktopSessionBridge } from '../../renderer/workhub-session-port.js'; -test('WorkHub candidates follow the resolved Coordination Session Host only', async () => { +test('WorkHub projections follow the resolved Coordination Session Host only', async () => { const sessionA = desktopSessionKey({ hostId: 'host-a', sessionId: 'ordinary-a' }); const sessionB = desktopSessionKey({ hostId: 'host-b', sessionId: 'ordinary-b' }); let coordinationSessionId: string | undefined; @@ -44,19 +41,9 @@ test('WorkHub candidates follow the resolved Coordination Session Host only', as completeHostIds: ['host-a', 'host-b'], }), listTurns: async () => [], - create: async () => { - throw new Error('unscoped create must not be used'); - }, - send: async () => ({ ok: true, turnId: 'turn' }), - stop: async () => undefined, + queryMessageExecutions: async () => ({ resolutions: [] }), subscribeChanges: () => () => undefined, }; - const createdFor: string[] = []; - const createOnCoordinationHost = async (coordinationId: string) => { - createdFor.push(coordinationId); - const { hostId } = parseDesktopSessionKey(coordinationId); - return ordinarySession(hostId === 'host-a' ? sessionA : sessionB); - }; const scopeSessions = () => { const generation = coordinationGeneration; return scopeWorkHubSessionsToCoordinationHost( @@ -65,7 +52,6 @@ test('WorkHub candidates follow the resolved Coordination Session Host only', as sessionId: coordinationSessionId, isCurrent: () => generation === coordinationGeneration, }, - createOnCoordinationHost, ); }; let sessions = scopeSessions(); @@ -93,13 +79,11 @@ test('WorkHub candidates follow the resolved Coordination Session Host only', as }); const unresolvedList = sessions.list(); - const unresolvedCreate = assert.rejects(sessions.create({ name: 'unsafe' }), /unresolved/); assert.deepEqual(await unresolvedList, []); - await unresolvedCreate; await Promise.resolve(); assert.deepEqual((await sessions.list()).map((session) => session.id), [sessionA]); await assert.rejects( - () => sessions.send(sessionB, { type: 'send', turnId: 'turn', text: 'wrong Host' }), + () => sessions.listTurns(sessionB), /another Runtime Host/, ); assert.deepEqual(await sessions.listWithCoverage?.(), { @@ -113,12 +97,7 @@ test('WorkHub candidates follow the resolved Coordination Session Host only', as assert.deepEqual(await sessions.list(), []); await Promise.resolve(); assert.deepEqual((await sessions.list()).map((session) => session.id), [sessionB]); - await assert.rejects(staleHostAScope.create({ name: 'stale A' }), /scope is revoked/); - assert.equal((await sessions.create({ name: 'current B' })).id, sessionB); - assert.deepEqual( - createdFor.map((sessionId) => parseDesktopSessionKey(sessionId).hostId), - ['host-b'], - ); + await assert.rejects(staleHostAScope.listTurns(sessionA), /scope is revoked/); stop(); }); diff --git a/apps/desktop/src/main/__tests__/workhub-session-port.test.ts b/apps/desktop/src/main/__tests__/workhub-session-port.test.ts index 8f3eb7d0c2..fbf60f1c43 100644 --- a/apps/desktop/src/main/__tests__/workhub-session-port.test.ts +++ b/apps/desktop/src/main/__tests__/workhub-session-port.test.ts @@ -31,7 +31,6 @@ import { createDesktopWorkHubCoordinationPort, projectWorkHubCoordinationTurns, } from '../../renderer/workhub-coordination-port.js'; -import { WorkHubSessionSubmitError } from '../../renderer/workhub-controller.js'; function desktopSession( id: string, @@ -56,6 +55,14 @@ const unusedTranscripts = { }, }; +const noMessageExecutions = async () => ({ + resolutions: [] as Array< + | { messageId: string; state: 'pending' } + | { messageId: string; state: 'cancelled' } + | { messageId: string; state: 'owned'; turnId: string; runId: string } + >, +}); + function transcriptsWith(messages: readonly StoredMessage[]) { return { open: async (sessionId: string, handler: (batch: DesktopTranscriptBatch) => void) => { @@ -148,8 +155,12 @@ test('projects the durable Coordination transcript into the WorkHub conversation text: 'Continue payments', state: 'completed', assignment: { + delegationId: 'payments-delegation', targetSessionId: 'payments', targetSessionName: 'Payments', + targetMessageId: 'payments-message', + targetTurnId: 'payments-turn', + feedbackState: 'accepted', }, updatedAt: 20, }]); @@ -189,7 +200,6 @@ test('Coordination transcript adapter emits an initial empty ready snapshot and }; }, }, - answer: async (input) => ({ turnId: input.turnId }), record: async (input) => ({ turnId: input.turnId }), candidates: async () => ({ candidateSetId: `sha256:${'a'.repeat(64)}`, @@ -291,6 +301,7 @@ test('desktop adapter rebuilds recent turns from the Session transcript and clos sessions: { list: async () => [], listTurns: async () => [], + queryMessageExecutions: noMessageExecutions, create: async () => { throw new Error('not used'); }, @@ -342,7 +353,6 @@ test('desktop adapter rebuilds recent turns from the Session transcript and clos }, }, projectName: () => 'Maka', - newTurnId: () => 'unused', }); assert.deepEqual(await adapter.recentTurns([{ sessionId }]), [{ @@ -366,6 +376,7 @@ test('desktop adapter cancels an unavailable transcript without hiding ready Ses sessions: { list: async () => [], listTurns: async () => [], + queryMessageExecutions: noMessageExecutions, create: async () => { throw new Error('not used'); }, send: async () => { throw new Error('not used'); }, stop: async () => {}, @@ -403,7 +414,6 @@ test('desktop adapter cancels an unavailable transcript without hiding ready Ses }, }, projectName: () => 'Maka', - newTurnId: () => 'unused', }); const turns = adapter.recentTurns([ @@ -448,6 +458,7 @@ test('desktop adapter projects Session catalog facts without owning copies', asy sessions: { list: async () => source, listTurns: async () => [], + queryMessageExecutions: noMessageExecutions, create: async () => { throw new Error('not used'); }, @@ -458,7 +469,6 @@ test('desktop adapter projects Session catalog facts without owning copies', asy subscribeChanges: () => () => {}, }, projectName: (projectId) => projectId === 'project-maka' ? 'Maka' : undefined, - newTurnId: () => 'unused', }); assert.deepEqual(await adapter.list(), [ @@ -506,266 +516,163 @@ test('desktop adapter projects Session catalog facts without owning copies', asy ]); }); -test('desktop adapter preserves per-Host catalog coverage for ownership reconciliation', async () => { - const localSessionId = desktopSessionKey({ hostId: 'local-host', sessionId: 'local' }); +test('desktop adapter rebuilds delegation feedback from the Message-owned execution Turn', async () => { + const sessions = [ + desktopSession('accepted'), + desktopSession('running', { status: 'running', runningTurnIds: ['turn-running'] }), + desktopSession('stale-running'), + desktopSession('recorded-running-only', { runningTurnIds: undefined }), + desktopSession('waiting', { + status: 'waiting_for_user', + runningTurnIds: ['turn-waiting'], + }), + desktopSession('completed', { + status: 'waiting_for_user', + runningTurnIds: ['later-turn'], + }), + desktopSession('failed'), + desktopSession('aborted'), + desktopSession('cancelled'), + desktopSession('recovering'), + ]; + const turns = new Map>([ + ['running', [{ turnId: 'turn-running', status: 'running', statusSource: 'recorded' }]], + ['stale-running', [{ + turnId: 'turn-stale-running', + status: 'running', + statusSource: 'recorded', + }]], + ['recorded-running-only', [{ + turnId: 'turn-recorded-running-only', + status: 'running', + statusSource: 'recorded', + }]], + ['waiting', [{ turnId: 'turn-waiting', status: 'running', statusSource: 'recorded' }]], + ['completed', [{ turnId: 'turn-completed', status: 'completed', statusSource: 'recorded' }]], + ['failed', [{ turnId: 'turn-failed', status: 'failed', statusSource: 'recorded' }]], + ['aborted', [{ turnId: 'turn-aborted', status: 'aborted', statusSource: 'recorded' }]], + ]); const adapter = createDesktopWorkHubSessionPort({ transcripts: unusedTranscripts, sessions: { - list: async () => [], - listWithCoverage: async () => ({ - sessions: [desktopSession(localSessionId)], - completeHostIds: ['local-host'], + list: async () => sessions, + listTurns: async (sessionId) => { + if (sessionId === 'recovering') throw new Error('Host is recovering'); + return turns.get(sessionId) ?? []; + }, + queryMessageExecutions: async (sessionId, messageIds) => ({ + resolutions: sessionId === 'accepted' + ? messageIds.map((messageId) => ({ messageId, state: 'pending' as const })) + : sessionId === 'cancelled' + ? messageIds.map((messageId) => ({ messageId, state: 'cancelled' as const })) + : sessionId === 'recovering' + ? [] + : messageIds.map((messageId) => ({ + messageId, + state: 'owned' as const, + turnId: `turn-${sessionId}`, + runId: `run-${sessionId}`, + })), }), - listTurns: async () => [], create: async () => { throw new Error('not used'); }, send: async () => { throw new Error('not used'); }, stop: async () => {}, subscribeChanges: () => () => {}, }, projectName: () => 'Maka', - newTurnId: () => 'unused', }); - - const catalog = await adapter.listCatalog?.(); - assert.ok(catalog); - assert.equal(catalog.sessions[0]?.target.sessionId, localSessionId); - assert.equal(catalog.isCompleteFor({ sessionId: localSessionId }), true); - assert.equal(catalog.isCompleteFor({ - sessionId: desktopSessionKey({ hostId: 'remote-host', sessionId: 'remote' }), - }), false); + const references = [ + ['accepted', 'turn-accepted'], + ['running', 'turn-running'], + ['stale-running', 'turn-stale-running'], + ['recorded-running-only', 'turn-recorded-running-only'], + ['waiting', 'turn-waiting'], + ['completed', 'turn-completed'], + ['failed', 'turn-failed'], + ['aborted', 'turn-aborted'], + ['cancelled', 'turn-cancelled'], + ['recovering', 'turn-recovering'], + ].map(([targetSessionId, targetTurnId]) => ({ + delegationId: `delegation-${targetSessionId}`, + targetSessionId: targetSessionId!, + targetMessageId: `message-${targetSessionId}`, + targetTurnId: targetTurnId!, + })); + + const feedback = await adapter.delegationFeedback(references); + + assert.deepEqual(feedback.map(({ delegationId, state }) => ({ delegationId, state })), [ + { delegationId: 'delegation-accepted', state: 'accepted' }, + { delegationId: 'delegation-running', state: 'running' }, + { delegationId: 'delegation-stale-running', state: 'accepted' }, + { delegationId: 'delegation-recorded-running-only', state: 'running' }, + { delegationId: 'delegation-waiting', state: 'waiting_for_user' }, + { delegationId: 'delegation-completed', state: 'completed' }, + { delegationId: 'delegation-failed', state: 'failed' }, + { delegationId: 'delegation-aborted', state: 'aborted' }, + { delegationId: 'delegation-cancelled', state: 'aborted' }, + { delegationId: 'delegation-recovering', state: 'recovering' }, + ]); }); -test('desktop adapter delegates create, send, and invalidation to Session APIs', async () => { - const calls: unknown[] = []; - let onChanged: (() => void) | undefined; +test('desktop adapter follows a delegated Message into its successor Turn', async () => { + const targetSessionId = desktopSessionKey({ hostId: 'local-host', sessionId: 'payments' }); const adapter = createDesktopWorkHubSessionPort({ - transcripts: unusedTranscripts, + transcripts: transcriptsWith([{ + type: 'user', + id: 'payment-message', + turnId: 'successor-turn', + ts: 2, + text: 'Continue payment recovery', + steeringEventId: 'payment-message', + }]), sessions: { - list: async () => [desktopSession('created', { + list: async () => [desktopSession(targetSessionId, { status: 'running', - runningTurnIds: ['turn-new'], + runningTurnIds: ['successor-turn'], })], - listTurns: async () => [], - create: async (input) => { - calls.push(['create', input]); - return desktopSession('created', { name: input.name }); - }, - send: async (sessionId, command) => { - calls.push(['send', sessionId, command]); - return { ok: true, turnId: command.turnId }; - }, - stop: async (sessionId, input) => { - calls.push(['stop', sessionId, input]); - }, - subscribeChanges: (handler) => { - onChanged = handler; - return () => calls.push(['unsubscribe']); - }, - }, - projectName: () => 'Maka', - newTurnId: () => 'turn-new', - }); - - const created = await adapter.create({ name: '实现导出发票 PDF 功能' }); - const turnId = adapter.reserveTurnId(); - const turn = await adapter.submit(created.target, '实现导出发票 PDF 功能', turnId); - await adapter.stop(created.target, 'turn-new'); - let invalidations = 0; - const unsubscribe = adapter.subscribe(() => { - invalidations += 1; - }); - onChanged?.(); - unsubscribe(); - - assert.equal(created.kind, 'ordinary'); - assert.deepEqual(turn, { turnId: 'turn-new' }); - assert.equal(invalidations, 1); - assert.deepEqual(calls, [ - ['create', { name: '实现导出发票 PDF 功能' }], - ['send', 'created', { type: 'send', turnId: 'turn-new', text: '实现导出发票 PDF 功能' }], - ['stop', 'created', { source: 'stop_button', expectedTurnId: 'turn-new' }], - ['unsubscribe'], - ]); -}); - -test('desktop adapter preserves when Session delivery steered an existing root Turn', async () => { - const adapter = createDesktopWorkHubSessionPort({ - transcripts: unusedTranscripts, - sessions: { - list: async () => [], - listTurns: async () => [], - create: async () => { - throw new Error('not used'); - }, - send: async (_sessionId, command) => ({ - ok: true, - turnId: command.turnId, - steered: true, + listTurns: async () => [ + { + turnId: 'admission-turn', + status: 'completed', + statusSource: 'recorded', + }, + { + turnId: 'successor-turn', + status: 'running', + statusSource: 'recorded', + }, + ], + queryMessageExecutions: async (_sessionId, messageIds) => ({ + resolutions: messageIds.map((messageId) => ({ + messageId, + state: 'owned' as const, + turnId: 'successor-turn', + runId: 'successor-run', + })), }), - stop: async () => {}, - subscribeChanges: () => () => {}, - }, - projectName: () => 'Maka', - newTurnId: () => 'turn-steered', - }); - - assert.deepEqual( - await adapter.submit( - { sessionId: 'busy' }, - '补充已有执行流', - adapter.reserveTurnId(), - ), - { turnId: 'turn-steered', steered: true }, - ); -}); - -test('desktop adapter distinguishes definite rejection from an unknown delivery outcome', async () => { - let outcome: 'throw' | 'unknown' | 'reject' = 'throw'; - const adapter = createDesktopWorkHubSessionPort({ - transcripts: unusedTranscripts, - sessions: { - list: async () => [], - listTurns: async () => [], create: async () => { throw new Error('not used'); }, - send: async () => { - if (outcome === 'throw') throw new Error('transport disconnected'); - if (outcome === 'unknown') { - return { - ok: false as const, - reason: 'outcome_unknown' as const, - messageId: 'reserved-turn', - skillInvocation: { loaded: [], failed: [], receipts: [] }, - }; - } - return { - ok: false as const, - reason: 'skill_invocation_failed' as const, - skillInvocation: { - loaded: [], - failed: [{ request: 'missing', reason: 'not_found' as const }], - receipts: [], - }, - }; - }, + send: async () => { throw new Error('not used'); }, stop: async () => {}, subscribeChanges: () => () => {}, }, projectName: () => 'Maka', - newTurnId: () => 'reserved-turn', - }); - - await assert.rejects( - adapter.submit({ sessionId: 'payment' }, '继续支付', 'reserved-turn'), - (error) => error instanceof WorkHubSessionSubmitError && error.admission === 'unknown', - ); - // The Host declining to prove the outcome must stay reconcilable; only a - // Host-owned refusal releases the reserved root. - outcome = 'unknown'; - await assert.rejects( - adapter.submit({ sessionId: 'payment' }, '继续支付', 'reserved-turn'), - (error) => error instanceof WorkHubSessionSubmitError && error.admission === 'unknown', - ); - outcome = 'reject'; - await assert.rejects( - adapter.submit({ sessionId: 'payment' }, '继续支付', 'reserved-turn'), - (error) => error instanceof WorkHubSessionSubmitError && error.admission === 'rejected', - ); -}); - -test('desktop adapter reconciles lost replies from authoritative transcript identity', async () => { - const cases: Array<{ - name: string; - message: StoredMessage; - expected: { kind: 'root'; turnId: string } | { kind: 'steered' } | { kind: 'unknown' }; - }> = [ - { - name: 'direct root', - message: { - type: 'user', id: 'user-root', turnId: 'reserved-turn', ts: 1, text: '开始支付', - }, - expected: { kind: 'root', turnId: 'reserved-turn' }, - }, - { - name: 'busy-race root', - message: { - type: 'user', id: 'reserved-turn', turnId: 'host-root', ts: 1, text: '开始支付', - }, - expected: { kind: 'root', turnId: 'host-root' }, - }, - { - name: 'steering', - message: { - type: 'user', - id: 'reserved-turn', - turnId: 'pre-existing-root', - steeringEventId: 'steering-event', - ts: 1, - text: '补充支付测试', - }, - expected: { kind: 'steered' }, - }, - { - name: 'unrelated message', - message: { - type: 'user', id: 'other-message', turnId: 'other-root', ts: 1, text: '其他工作', - }, - expected: { kind: 'unknown' }, - }, - ]; - - for (const fixture of cases) { - const adapter = createDesktopWorkHubSessionPort({ - transcripts: transcriptsWith([fixture.message]), - sessions: { - list: async () => [], - listTurns: async () => [], - create: async () => { throw new Error('not used'); }, - send: async () => { throw new Error('not used'); }, - stop: async () => {}, - subscribeChanges: () => () => {}, - }, - projectName: () => 'Maka', - newTurnId: () => 'reserved-turn', - }); - - assert.deepEqual( - await adapter.reconcileSubmission({ - sessionId: desktopSessionKey({ hostId: 'local-host', sessionId: fixture.name }), - }, 'reserved-turn'), - fixture.expected, - fixture.name, - ); - } -}); - -test('desktop adapter binds stop to the root Turn owned by the WorkHub submission', async () => { - const stopped: unknown[] = []; - const adapter = createDesktopWorkHubSessionPort({ - transcripts: unusedTranscripts, - sessions: { - list: async () => [], - listTurns: async () => [], - create: async () => { - throw new Error('not used'); - }, - send: async () => { - throw new Error('not used'); - }, - stop: async (sessionId, input) => { - stopped.push([sessionId, input]); - }, - subscribeChanges: () => () => {}, - }, - projectName: () => 'Maka', - newTurnId: () => 'unused', }); - await adapter.stop({ sessionId: 'payment' }, 'turn-workhub'); - - assert.deepEqual(stopped, [[ - 'payment', - { source: 'stop_button', expectedTurnId: 'turn-workhub' }, - ]]); + const references = [{ + delegationId: 'payment-delegation', + targetSessionId, + targetTurnId: 'admission-turn', + targetMessageId: 'payment-message', + }]; + assert.deepEqual(await adapter.delegationFeedback(references), [{ + delegationId: 'payment-delegation', + state: 'running', + }]); }); test('desktop adapter derives stable origin evidence from the existing Session log', async () => { @@ -782,6 +689,7 @@ test('desktop adapter derives stable origin evidence from the existing Session l { userPromptPreview: '把风险按高、中、低分组' }, ]; }, + queryMessageExecutions: noMessageExecutions, create: async () => { throw new Error('not used'); }, @@ -792,7 +700,6 @@ test('desktop adapter derives stable origin evidence from the existing Session l subscribeChanges: () => () => {}, }, projectName: () => 'Maka', - newTurnId: () => 'unused', }); const first = await adapter.routingEvidence([{ sessionId: 'payment' }]); diff --git a/apps/desktop/src/main/__tests__/workhub-surface-flow.test.ts b/apps/desktop/src/main/__tests__/workhub-surface-flow.test.ts index b6f1eb03ec..98af7045e5 100644 --- a/apps/desktop/src/main/__tests__/workhub-surface-flow.test.ts +++ b/apps/desktop/src/main/__tests__/workhub-surface-flow.test.ts @@ -24,6 +24,7 @@ import { renderToStaticMarkup } from 'react-dom/server'; import { AstryxLocaleProvider, LocaleProvider } from '@maka/ui'; import { WorkHubCoordinationStatus, + WorkHubCoordinationTurnView, WorkHubProjectionRefreshGate, WorkHubSurfaceRouteGate, submitAndRecordWorkHubSurfaceInput, @@ -34,9 +35,11 @@ import { workHubSubmissionClearsDraft, } from '../../renderer/workhub-surface.js'; import { - createLegacyWorkHubControllerForTests as createWorkHubController, + createWorkHubController, WORKHUB_ROUTING_STRATEGY_ID, type WorkHubController, + type WorkHubCoordinationTurn, + type WorkHubDelegationExecutionState, type WorkHubSubmitInput, } from '../../renderer/workhub-controller.js'; import { WorkHubSendLease } from '../../renderer/workhub-send-lease.js'; @@ -107,6 +110,52 @@ test('Coordination lifecycle keeps a visible loading state and exposes failure r assert.match(failed, />Retry { + const states: Array<[WorkHubDelegationExecutionState, string]> = [ + ['accepted', 'Accepted'], + ['running', 'Running'], + ['waiting_for_user', 'Waiting for you'], + ['completed', 'Completed'], + ['failed', 'Failed'], + ['aborted', 'Aborted'], + ['recovering', 'Recovering'], + ]; + for (const [state, label] of states) { + const turn: WorkHubCoordinationTurn = { + messageId: 'assignment-1', + turnId: 'action-1', + text: 'Continue payments', + state: 'completed', + assignment: { + delegationId: 'delegation-1', + targetSessionId: 'payment', + targetSessionName: 'Payments', + targetMessageId: 'payment-message', + targetTurnId: 'payment-turn', + feedbackState: state, + }, + updatedAt: 10, + }; + const markup = renderToStaticMarkup( + createElement(LocaleProvider, { + locale: 'en', + children: createElement(AstryxLocaleProvider, { + children: createElement(WorkHubCoordinationTurnView, { + turn, + projection: { sessions: [], turns: [] }, + locale: 'en', + onOpenSession: () => undefined, + }), + }), + }), + ); + assert.match(markup, / @@ -747,6 +772,15 @@ function workHubCopy(locale: UiLocale) { delivery_failed: '输入未能送达,请重试。', }, scrollToBottom: '滚动到底部', archived: '已归档', states: { active: '活跃', running: '进行中', waiting_for_user: '等待你', blocked: '受阻', aborted: '已中止' }, + delegationStates: { + accepted: '已接收', + running: '进行中', + waiting_for_user: '等待你', + completed: '已完成', + failed: '失败', + aborted: '已中止', + recovering: '正在恢复', + }, turnStates: { running: '进行中', completed: '已完成', aborted: '已中止', failed: '失败' }, } as const; } @@ -780,6 +814,15 @@ function workHubCopy(locale: UiLocale) { delivery_failed: 'The input could not be delivered. Try again.', }, scrollToBottom: 'Scroll to bottom', archived: 'Archived', states: { active: 'Active', running: 'Running', waiting_for_user: 'Waiting for you', blocked: 'Blocked', aborted: 'Aborted' }, + delegationStates: { + accepted: 'Accepted', + running: 'Running', + waiting_for_user: 'Waiting for you', + completed: 'Completed', + failed: 'Failed', + aborted: 'Aborted', + recovering: 'Recovering', + }, turnStates: { running: 'Running', completed: 'Completed', aborted: 'Aborted', failed: 'Failed' }, } as const; } diff --git a/docs/architecture/workhub-coordination-session-adr.md b/docs/architecture/workhub-coordination-session-adr.md index eade3fbfb9..7b49723565 100644 --- a/docs/architecture/workhub-coordination-session-adr.md +++ b/docs/architecture/workhub-coordination-session-adr.md @@ -98,6 +98,7 @@ transcripts, such as: delegationId coordinationTurnId targetSessionId +targetMessageId targetTurnId disposition ``` @@ -131,6 +132,22 @@ so WorkHub does not own a second recovery state machine or compensation chain. The `delegation_assigned` record itself projects the visible WorkHub turn; the renderer does not append a second summary. +The first-response contract is hybrid. The atomic `delegation_assigned` record is +an immediate durable acknowledgement, so WorkHub confirms acceptance without +waiting for target execution. The target Message is the stable delegation +identity; `targetTurnId` records only its admission location. WorkHub asks the +target Message authority which Turn durably consumed or admitted that Message, +then joins the resolved Turn's recorded lifecycle and the target Session's exact +live-Turn membership to project `running`, `waiting_for_user`, `completed`, +`failed`, and `aborted`. This remains correct when an unconsumed steering Message +is folded into a successor Turn or recovery aggregates several pending Messages +under one new Turn. A durable cancellation tombstone for a retracted queued +Message resolves the delegation to `aborted`. If the target authority is +temporarily unreadable, WorkHub projects `recovering` rather than inventing a +terminal result. These execution states are never appended as mutable Coordination +records; Session change notifications invalidate the projection and opening +WorkHub after restart rebuilds it from the same link and target facts. + The renderer persists only a Host-scoped action id until acknowledgement. Composer draft text uses a separate storage key and lifecycle. A reload therefore preserves idempotency without freezing old text or coupling draft edits to Host authority. @@ -156,7 +173,8 @@ been committed. - Coordination Session role representation, lazy creation, durable lookup, recovery, per-Host UI resolution, persistent transcript, closed dispositions, and the Action Gate are implemented. Durable delegation linkage is encoded in - that transcript; target lifecycle projection, linked correction, and destructive + that transcript; target lifecycle projection and the hybrid first-response + contract are implemented as rebuildable reads. Linked correction and destructive replacement/Stop recovery remain later work. Reevaluate the per-Host decision if supported workflows require one WorkHub diff --git a/packages/runtime-host/src/__tests__/message-coordinator.test.ts b/packages/runtime-host/src/__tests__/message-coordinator.test.ts index 48f24ad74a..9b30146285 100644 --- a/packages/runtime-host/src/__tests__/message-coordinator.test.ts +++ b/packages/runtime-host/src/__tests__/message-coordinator.test.ts @@ -130,6 +130,58 @@ test('consumes an atomically committed active-target admission exactly once', as assert.equal(fixture.drainRequests(), 0); }); +test('idle recovery resolves differently preassigned Messages to their shared successor Turn', async () => { + const fixture = createFixture(); + fixture.setRootState({ kind: 'idle' }); + for (const [messageId, turnId, runId] of [ + ['workhub-message-a', 'preassigned-turn-a', 'preassigned-run-a'], + ['workhub-message-b', 'preassigned-turn-b', 'preassigned-run-b'], + ] as const) { + const content = { text: `recover ${messageId}` }; + await fixture.admissions.commitMessageAdmission({ + sessionId: ROOT.sessionId, + turnId, + runId, + messageId, + content, + submittedContentDigest: messageContentDigest(content), + submittedPlacement: 'current_turn', + placement: 'current_turn', + disposition: 'steering', + admittedAt: 10, + }); + } + + await fixture.coordinator.consumePendingAdmissions([ROOT.sessionId]); + const resolved = await fixture.coordinator.handlers['turn.message.execution.query']( + { + sessionId: ROOT.sessionId, + messageIds: ['workhub-message-a', 'workhub-message-b'], + }, + operationContext(), + ); + + assert.deepEqual(resolved, { + ok: true, + result: { + resolutions: [ + { + messageId: 'workhub-message-a', + state: 'owned', + turnId: 'recovered-turn', + runId: 'durable-run', + }, + { + messageId: 'workhub-message-b', + state: 'owned', + turnId: 'recovered-turn', + runId: 'durable-run', + }, + ], + }, + }); +}); + test('idle submit starts exactly one root Turn and retry identity is connection-independent', async () => { const fixture = createFixture(); fixture.setRootState({ kind: 'idle' }); @@ -202,27 +254,75 @@ test('message query reports only durable cancellation proof', async () => { await submit(fixture, 'cancelled-message', 'discard me', 'next_turn'); await submit(fixture, 'accepted-message', 'waiting', 'next_turn'); await fixture.coordinator.cancelMessages(ROOT.sessionId, ['cancelled-message']); + const result = await fixture.coordinator.handlers['turn.message.query']( + { + sessionId: ROOT.sessionId, + messageIds: ['cancelled-message', 'accepted-message', 'unknown-message'], + }, + operationContext(), + ); + + assert.deepEqual(result, { + ok: true, + result: { cancelledMessageIds: ['cancelled-message'] }, + }); +}); + +test('message execution query reports the Turn that durably owns each Message', async () => { + const fixture = createFixture(); + const pendingContent = { text: 'not handed off yet' }; + await fixture.admissions.commitMessageAdmission({ + ...ROOT, + messageId: 'pending-message', + content: pendingContent, + submittedContentDigest: messageContentDigest(pendingContent), + submittedPlacement: 'current_turn', + placement: 'current_turn', + disposition: 'steering', + admittedAt: 10, + }); fixture.receipts.set( 'handed-off-message', - sourceReceipt('handed-off-message', 'delivered', 'current_turn', 'steering'), + sourceReceipt( + 'handed-off-message', + 'delivered by successor', + 'current_turn', + 'steering', + 'successor-turn', + ), ); + fixture.events.push(steeringEvent('steered-message', 'consumed by admission Turn')); - const result = await fixture.coordinator.handlers['turn.message.query']( + const result = await fixture.coordinator.handlers['turn.message.execution.query']( { sessionId: ROOT.sessionId, - messageIds: [ - 'cancelled-message', - 'accepted-message', - 'handed-off-message', - 'unknown-message', - ], + messageIds: ['pending-message', 'handed-off-message', 'steered-message', 'unknown-message'], }, operationContext(), ); assert.deepEqual(result, { ok: true, - result: { cancelledMessageIds: ['cancelled-message'] }, + result: { + resolutions: [ + { + messageId: 'pending-message', + state: 'pending', + }, + { + messageId: 'handed-off-message', + state: 'owned', + turnId: 'successor-turn', + runId: 'durable-run', + }, + { + messageId: 'steered-message', + state: 'owned', + turnId: ROOT.turnId, + runId: ROOT.runId, + }, + ], + }, }); }); @@ -743,6 +843,18 @@ test('entry retract removes one queued entry, replays its outcome, and rejects s fixture.coordinator.projection(ROOT.sessionId).steering.map((entry) => entry.messageId), ['steer-1'], ); + assert.deepEqual( + await fixture.coordinator.handlers['turn.message.execution.query']( + { sessionId: ROOT.sessionId, messageIds: ['follow-1'] }, + operationContext(), + ), + { + ok: true, + result: { + resolutions: [{ messageId: 'follow-1', state: 'cancelled' }], + }, + }, + ); const retry = await fixture.coordinator.handlers['queue.entry.retract']( { diff --git a/packages/runtime-host/src/__tests__/protocol.test.ts b/packages/runtime-host/src/__tests__/protocol.test.ts index 8c4b8e3f2b..1c5241a5e9 100644 --- a/packages/runtime-host/src/__tests__/protocol.test.ts +++ b/packages/runtime-host/src/__tests__/protocol.test.ts @@ -333,6 +333,10 @@ describe('Runtime Host bootstrap protocol', () => { assert.ok(RUNTIME_HOST_COMPATIBILITY_EPOCH > 50); }); + test('publishes a new compatibility epoch for Message execution ownership', () => { + assert.ok(RUNTIME_HOST_COMPATIBILITY_EPOCH > 61); + }); + test('publishes a new compatibility epoch for exact Session Connection identity', () => { assert.ok(RUNTIME_HOST_COMPATIBILITY_EPOCH > 56); }); @@ -1185,9 +1189,14 @@ describe('Runtime Host bootstrap protocol', () => { operation: 'turn.message.query' as const, input: { sessionId: 'session-1', - messageIds: ['message-1', 'message-2'], + messageIds: ['message-1', 'message-2', 'message-3'], }, }; + const executionQuery = { + requestId: 'execution-query-request-1', + operation: 'turn.message.execution.query' as const, + input: query.input, + }; const submit = { requestId: 'submit-request-1', operation: 'turn.message.submit' as const, @@ -1216,6 +1225,36 @@ describe('Runtime Host bootstrap protocol', () => { }, }; assert.deepEqual(decodeClientFrame(query), query); + assert.deepEqual(decodeClientFrame(executionQuery), executionQuery); + const queried = { + requestId: executionQuery.requestId, + operation: executionQuery.operation, + ok: true as const, + result: { + resolutions: [ + { messageId: 'message-1', state: 'pending' as const }, + { + messageId: 'message-2', + state: 'owned' as const, + turnId: 'turn-2', + runId: 'run-2', + }, + { messageId: 'message-3', state: 'cancelled' as const }, + ], + }, + }; + assert.deepEqual(decodeHostFrame(queried), queried); + assert.throws( + () => + decodeHostFrame({ + ...queried, + result: { + ...queried.result, + resolutions: [...queried.result.resolutions, ...queried.result.resolutions], + }, + }), + isInvalidFrame, + ); assert.deepEqual(decodeClientFrame(submit), submit); assert.deepEqual(decodeClientFrame(retract), retract); assert.deepEqual(decodeClientFrame(interrupt), interrupt); diff --git a/packages/runtime-host/src/__tests__/root-turn-coordinator.test.ts b/packages/runtime-host/src/__tests__/root-turn-coordinator.test.ts index 08eb929d6c..a0f2f62842 100644 --- a/packages/runtime-host/src/__tests__/root-turn-coordinator.test.ts +++ b/packages/runtime-host/src/__tests__/root-turn-coordinator.test.ts @@ -55,7 +55,7 @@ import type { BackendCompactHistoryInput, BackendSendInput, } from '@maka/core/backend-types'; -import type { SessionEvent } from '@maka/core/events'; +import { messageContentDigest, type SessionEvent } from '@maka/core/events'; import { WORKHUB_COORDINATION_SESSION_ID, WORKHUB_COORDINATION_SESSION_ROLE, diff --git a/packages/runtime-host/src/protocol/index.ts b/packages/runtime-host/src/protocol/index.ts index 253c42fb20..a359bcdefd 100644 --- a/packages/runtime-host/src/protocol/index.ts +++ b/packages/runtime-host/src/protocol/index.ts @@ -94,7 +94,9 @@ export const RUNTIME_HOST_REGISTRATION_SCHEMA_VERSION = 1 as const; export const RUNTIME_HOST_PROTOCOL_VERSION = 0 as const; // Increment when the same protocol version no longer guarantees safe Client-Host // interoperability. Mismatches are rejected before domain commands are admitted. -export const RUNTIME_HOST_COMPATIBILITY_EPOCH = 66 as const; +export const RUNTIME_HOST_COMPATIBILITY_EPOCH = 67 as const; +// 67: Message lifecycle queries expose durable execution ownership and +// cancellation. Older peers cannot decode or provide the closed proof list. // 66: Peer Mesh queries expose one canonical transit selection and runtime metrics. // 65: live `tool_start` frames may carry optional `intent` / `argsPreview` // keys. Older Clients decode the event with a strict allowed-key list and tear diff --git a/packages/runtime-host/src/protocol/message.ts b/packages/runtime-host/src/protocol/message.ts index ea4df35d27..c30546427c 100644 --- a/packages/runtime-host/src/protocol/message.ts +++ b/packages/runtime-host/src/protocol/message.ts @@ -123,6 +123,25 @@ export interface TurnMessageQueryResult { readonly cancelledMessageIds: readonly string[]; } +export interface TurnMessageExecutionQueryInput { + readonly sessionId: string; + readonly messageIds: readonly string[]; +} + +export interface TurnMessageExecutionQueryResult { + readonly resolutions: readonly TurnMessageExecutionResolution[]; +} + +export type TurnMessageExecutionResolution = + | { readonly messageId: string; readonly state: 'pending' } + | { readonly messageId: string; readonly state: 'cancelled' } + | { + readonly messageId: string; + readonly state: 'owned'; + readonly turnId: string; + readonly runId: string; + }; + export interface QueueRetractInput { readonly originHostEpoch: string; readonly sessionId: string; @@ -202,6 +221,13 @@ export const MESSAGE_OPERATION_SPECS = { decodeInput: decodeTurnMessageQueryInput, decodeOutput: decodeTurnMessageQueryResult, }), + 'turn.message.execution.query': defineOperation({ + mode: 'query', + availability: 'ready', + errors: MESSAGE_OPERATION_ERRORS, + decodeInput: decodeTurnMessageExecutionQueryInput, + decodeOutput: decodeTurnMessageExecutionQueryResult, + }), 'turn.message.submit': defineOperation({ mode: 'command', availability: 'ready', @@ -341,6 +367,59 @@ function decodeTurnMessageQueryResult(value: unknown): TurnMessageQueryResult { return { cancelledMessageIds }; } +function decodeTurnMessageExecutionQueryInput(value: unknown): TurnMessageExecutionQueryInput { + return decodeTurnMessageQueryInput(value); +} + +function decodeTurnMessageExecutionQueryResult(value: unknown): TurnMessageExecutionQueryResult { + const record = requireExactRecord(value, 'turn.message.execution.query result', ['resolutions']); + if (!Array.isArray(record.resolutions) || record.resolutions.length > MESSAGE_QUEUE_MAX_ENTRIES) { + throw invalidProtocolFrame('Invalid turn.message.execution.query resolutions'); + } + const resolutions = record.resolutions.map((value): TurnMessageExecutionResolution => { + const resolution = requireRecord(value, 'turn.message.execution.query resolution'); + if (resolution.state === 'pending') { + assertExactKeys(resolution, 'turn.message.execution.query pending resolution', [ + 'messageId', + 'state', + ]); + return { + messageId: requireEntityId(resolution.messageId, 'messageId'), + state: 'pending', + }; + } + if (resolution.state === 'cancelled') { + assertExactKeys(resolution, 'turn.message.execution.query cancelled resolution', [ + 'messageId', + 'state', + ]); + return { + messageId: requireEntityId(resolution.messageId, 'messageId'), + state: 'cancelled', + }; + } + if (resolution.state === 'owned') { + assertExactKeys(resolution, 'turn.message.execution.query owned resolution', [ + 'messageId', + 'state', + 'turnId', + 'runId', + ]); + return { + messageId: requireEntityId(resolution.messageId, 'messageId'), + state: 'owned', + turnId: requireEntityId(resolution.turnId, 'turnId'), + runId: requireEntityId(resolution.runId, 'runId'), + }; + } + throw invalidProtocolFrame('Invalid turn.message.execution.query resolution state'); + }); + if (new Set(resolutions.map(({ messageId }) => messageId)).size !== resolutions.length) { + throw invalidProtocolFrame('Duplicate turn.message.execution.query messageId'); + } + return { resolutions }; +} + function decodeTurnMessageSubmitResult(value: unknown): TurnMessageSubmitResult { const record = requireRecord(value, 'turn.message.submit result'); if (record.disposition === 'turn_started') { diff --git a/packages/runtime-host/src/protocol/operations.ts b/packages/runtime-host/src/protocol/operations.ts index d9176c2657..392347d4f7 100644 --- a/packages/runtime-host/src/protocol/operations.ts +++ b/packages/runtime-host/src/protocol/operations.ts @@ -317,6 +317,7 @@ export const REMOTE_OWNER_OPERATION_GRANTS = Object.freeze([ 'subscription.open', 'task.ledger.query', 'turn.interrupt', + 'turn.message.execution.query', 'turn.message.query', 'turn.message.submit', 'turn.query', diff --git a/packages/runtime-host/src/server/message-coordinator.ts b/packages/runtime-host/src/server/message-coordinator.ts index 1db629328d..0f3ec3f083 100644 --- a/packages/runtime-host/src/server/message-coordinator.ts +++ b/packages/runtime-host/src/server/message-coordinator.ts @@ -338,6 +338,7 @@ const HOST_EPOCH_PATTERN = /^[A-Za-z0-9_-]{1,128}$/u; export class HostMessageCoordinator implements RuntimeMessageAuthority { readonly handlers: MessageOperationHandlerMap = { 'turn.message.query': (input) => this.queryMessages(input), + 'turn.message.execution.query': (input) => this.queryMessageExecutions(input), 'turn.message.submit': (input, context) => this.submit(input, context), 'queue.retract': (input) => this.retract(input), 'queue.entry.retract': (input) => this.retractQueuedEntry(input), @@ -411,6 +412,71 @@ export class HostMessageCoordinator implements RuntimeMessageAuthority { return success({ cancelledMessageIds }); } + async queryMessageExecutions(input: { + sessionId: string; + messageIds: readonly string[]; + }): Promise< + MessageOutcome<{ + resolutions: Array< + | { messageId: string; state: 'pending' } + | { messageId: string; state: 'cancelled' } + | { messageId: string; state: 'owned'; turnId: string; runId: string } + >; + }> + > { + const resolutions: Array< + | { messageId: string; state: 'pending' } + | { messageId: string; state: 'cancelled' } + | { messageId: string; state: 'owned'; turnId: string; runId: string } + > = []; + for (const messageId of input.messageIds) { + const receipt = await this.#durableProof.readRootTurnSourceMessageReceipt( + input.sessionId, + messageId, + ); + if ( + receipt?.admission.sessionId === input.sessionId && + receipt.sourceMessage.messageId === messageId + ) { + // A root source receipt is the latest durable ownership proof and + // therefore outranks the steering location from which a Message may + // have been folded into this successor. + resolutions.push({ + messageId, + state: 'owned', + turnId: receipt.admission.turnId, + runId: receipt.admission.runId, + }); + continue; + } + const steering = await this.#durableProof.readImmutableSteeringMessageProof( + input.sessionId, + messageId, + ); + if ( + steering?.event.sessionId === input.sessionId && + steering.event.refs?.providerEventId === messageId + ) { + resolutions.push({ + messageId, + state: 'owned', + turnId: steering.event.turnId, + runId: steering.event.runId, + }); + continue; + } + if (await this.#admissions.hasCancelledMessageAdmission(input.sessionId, messageId)) { + resolutions.push({ messageId, state: 'cancelled' }); + continue; + } + const pending = await this.#admissions.readMessageAdmission(input.sessionId, messageId); + if (pending?.sessionId === input.sessionId && pending.messageId === messageId) { + resolutions.push({ messageId, state: 'pending' }); + } + } + return success({ resolutions }); + } + retireSessions(sessionIds: readonly string[]): void { for (const sessionId of new Set(sessionIds)) { const state = this.#sessions.get(sessionId); diff --git a/packages/runtime-host/src/server/operation-dispatcher.ts b/packages/runtime-host/src/server/operation-dispatcher.ts index 8208c5ad08..77af9e17a5 100644 --- a/packages/runtime-host/src/server/operation-dispatcher.ts +++ b/packages/runtime-host/src/server/operation-dispatcher.ts @@ -90,6 +90,7 @@ export type ConnectionEffectOperationKey = Extract< export type MessageOperationKey = Extract< OperationKey, | 'turn.message.query' + | 'turn.message.execution.query' | 'turn.message.submit' | 'queue.retract' | 'queue.entry.retract'