diff --git a/.maka-shots/workhub-context-after.png b/.maka-shots/workhub-context-after.png new file mode 100644 index 0000000000..8b82852c06 Binary files /dev/null and b/.maka-shots/workhub-context-after.png differ diff --git a/.maka-shots/workhub-context-before.png b/.maka-shots/workhub-context-before.png new file mode 100644 index 0000000000..a2903e2e52 Binary files /dev/null and b/.maka-shots/workhub-context-before.png differ diff --git a/apps/desktop/e2e/workhub-reconstruction.spec.ts b/apps/desktop/e2e/workhub-reconstruction.spec.ts index d81dec3861..e853d7f733 100644 --- a/apps/desktop/e2e/workhub-reconstruction.spec.ts +++ b/apps/desktop/e2e/workhub-reconstruction.spec.ts @@ -44,14 +44,16 @@ test('WorkHub rebuilds Session conversation after navigating away and back', asy }), ).toBeVisible(); - const routedPrompt = `继续${sessionName},补充重复投递测试点。`; + const routedPrompt = '继续这个工作,补充重复投递测试点。'; const workHubComposer = page.locator( '.workhub-surface .maka-composer-editor [contenteditable="true"]', ); await workHubComposer.fill(routedPrompt); await workHubComposer.press('Enter'); await expect(page.locator('.workhub-submitted').last()).toBeVisible(); - await page.locator('.workhub-submitted > button').last().click(); + await page.locator('.workhub-turn', { hasText: routedPrompt }) + .locator('.workhub-submitted > button') + .click(); await expect(page.getByRole('main', { name: 'WorkHub' })).toBeHidden(); await page.getByRole('button', { name: 'WorkHub', exact: true }).click(); @@ -62,3 +64,47 @@ test('WorkHub rebuilds Session conversation after navigating away and back', asy }), ).toBeVisible(); }); + +test('WorkHub handles a first natural-language correction', async ({ + window: page, +}) => { + const composer = page.locator(COMPOSER_INPUT); + await composer.fill('检查支付回调重复投递时的幂等性'); + await composer.press('Enter'); + await expect(page.getByRole('button', { name: '重新生成' })).toHaveCount(1, { + timeout: 20_000, + }); + await page.evaluate(async () => { + await window.maka.settings.updateClient({ workHub: { enabled: true } }); + }); + await expect(page.getByRole('main', { name: 'WorkHub' })).toBeVisible(); + await page.evaluate(async () => { + await window.maka.sessions.create({ name: '登录稳定性' }); + }); + await expect(page.getByText('2 项工作', { exact: true })).toBeVisible(); + + const workHubComposer = page.locator( + '.workhub-surface .maka-composer-editor [contenteditable="true"]', + ); + await workHubComposer.fill('继续这个工作,补充重复投递测试点。'); + await workHubComposer.press('Enter'); + const continuedTurn = page.locator('.workhub-turn', { + hasText: '继续这个工作,补充重复投递测试点。', + }); + await expect( + continuedTurn.locator('.workhub-submitted-session strong'), + ).toHaveText('检查支付回调重复投递时的幂等性'); + + await workHubComposer.fill('不是这个,换成登录稳定性,补充刷新令牌失败判定。'); + await expect( + page.locator('.workhub-surface').getByRole('button', { name: '发送' }), + ).toBeEnabled(); + await workHubComposer.press('Enter'); + + await expect(page.locator('.workhub-correction-note').last()).toBeVisible(); + await expect( + page.locator('.workhub-turn', { + hasText: '不是这个,换成登录稳定性,补充刷新令牌失败判定。', + }).locator('.workhub-submitted-session strong'), + ).toHaveText('登录稳定性'); +}); diff --git a/apps/desktop/src/main/__tests__/runtime-host-session-catalog-preload.test.ts b/apps/desktop/src/main/__tests__/runtime-host-session-catalog-preload.test.ts index b90e6bbd86..99441165b4 100644 --- a/apps/desktop/src/main/__tests__/runtime-host-session-catalog-preload.test.ts +++ b/apps/desktop/src/main/__tests__/runtime-host-session-catalog-preload.test.ts @@ -20,7 +20,10 @@ import assert from 'node:assert/strict'; import test from 'node:test'; import type { DesktopSessionSummary } from '../../preload/bridge-contract.js'; -import { collectRuntimeHostSessionCatalogs } from '../../preload/runtime-host-session-catalog.js'; +import { + collectRuntimeHostSessionCatalogs, + collectRuntimeHostSessionCatalogsWithCoverage, +} from '../../preload/runtime-host-session-catalog.js'; function session(id: string, activityAt: number): DesktopSessionSummary { return { id, activityAt } as DesktopSessionSummary; @@ -36,6 +39,16 @@ test('keeps healthy Host catalogs when another Host rejects', async () => { assert.deepEqual(sessions.map(({ id }) => id), ['newer', 'older']); }); +test('reports exactly which Host catalogs are complete', async () => { + const catalog = await collectRuntimeHostSessionCatalogsWithCoverage([ + { hostId: 'local', sessions: Promise.resolve([session('local-session', 1)]) }, + { hostId: 'remote', sessions: Promise.reject(new Error('remote unavailable')) }, + ]); + + assert.deepEqual(catalog.sessions.map(({ id }) => id), ['local-session']); + assert.deepEqual(catalog.completeHostIds, ['local']); +}); + test('fails when every Host catalog rejects', async () => { await assert.rejects( collectRuntimeHostSessionCatalogs([ 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 a59a82cf65..4514ea326b 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 @@ -609,7 +609,7 @@ test("queues a mid-turn send as steering when the Host reports the session busy" assert.deepEqual(submits, [ { sessionId: "session-1", - messageId: "id-1", + messageId: "turn-1", content: { text: "also check the tests", inlineReferences: [] }, placement: "current_turn", }, @@ -629,6 +629,7 @@ test("queues a mid-turn send as steering when the Host reports the session busy" test("starts the turn from the queued message when the busy race resolves idle", async () => { const changes: unknown[] = []; + const submits: unknown[] = []; const ipc = ipcHarness(); registerExecutionIpc( { @@ -641,10 +642,13 @@ test("starts the turn from the queued message when the busy race resolves idle", "Session already has an active root Turn", ); }, - submitMessage: async () => ({ - disposition: "turn_started", - turnId: "turn-9", - }), + submitMessage: async (input) => { + submits.push(input); + return { + disposition: "turn_started", + turnId: "turn-9", + }; + }, }), observer: unusedObserver(), attachmentApprovals: createAttachmentApprovalRegistry(), @@ -671,6 +675,12 @@ test("starts the turn from the queued message when the busy race resolves idle", inlineReferences: [], skillInvocation: { loaded: [], failed: [], receipts: [] }, }); + assert.deepEqual(submits, [{ + sessionId: "session-1", + messageId: "turn-1", + content: { text: "also check the tests", inlineReferences: [] }, + placement: "current_turn", + }]); assert.deepEqual(changes, [ { reason: "status-change", sessionId: "session-1", turnId: "turn-9" }, ]); diff --git a/apps/desktop/src/main/__tests__/workhub-controller.test.ts b/apps/desktop/src/main/__tests__/workhub-controller.test.ts index 509c131797..821c69c24f 100644 --- a/apps/desktop/src/main/__tests__/workhub-controller.test.ts +++ b/apps/desktop/src/main/__tests__/workhub-controller.test.ts @@ -18,16 +18,36 @@ */ import assert from 'node:assert/strict'; +import { existsSync, readFileSync } from 'node:fs'; import test from 'node:test'; import { createWorkHubController, + WorkHubSessionSubmitError, WORKHUB_ROUTING_STRATEGY_ID, type WorkHubSessionFacts, type WorkHubSessionPort, } from '../../renderer/workhub-controller.js'; -test('binds the controller to the immutable WH-R2.3 strategy ID', () => { - assert.equal(WORKHUB_ROUTING_STRATEGY_ID, 'wh-r2.3-session-core-evidence'); +const appShellUrl = [ + new URL('../../renderer/app-shell.tsx', import.meta.url), + new URL('../../../src/renderer/app-shell.tsx', import.meta.url), +].find((candidate) => existsSync(candidate)); + +if (!appShellUrl) throw new Error('Could not locate renderer/app-shell.tsx'); + +test('binds the controller to the immutable WH-R2.4 strategy ID', () => { + assert.equal(WORKHUB_ROUTING_STRATEGY_ID, 'wh-r2.4-session-context-continuity'); +}); + +test('binds the WorkHub controller to the app runtime rather than project refreshes', () => { + const source = readFileSync(appShellUrl, 'utf8'); + + assert.match(source, /workHubControllerRef\s*=\s*useRef/u); + assert.match(source, /workHubProjectsRef\.current\s*=\s*projects/u); + assert.doesNotMatch( + source, + /useMemo\(\(\)\s*=>\s*createWorkHubController\([\s\S]*?\),\s*\[projects\]\)/u, + ); }); function session( @@ -47,6 +67,7 @@ function session( } function port(sessions: WorkHubSessionFacts[]): WorkHubSessionPort { + let nextTurnId = 0; return { list: async () => sessions, recentTurns: async () => [], @@ -54,9 +75,11 @@ function port(sessions: WorkHubSessionFacts[]): WorkHubSessionPort { 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 () => {}, subscribe: () => () => {}, }; @@ -194,7 +217,7 @@ test('submit sends an explicitly targeted request to that Session', async () => assert.deepEqual(result, { kind: 'submitted', - strategyId: 'wh-r2.3-session-core-evidence', + strategyId: WORKHUB_ROUTING_STRATEGY_ID, requestId: 'request-1', target: { sessionId: 'payment' }, turnId: 'turn-payment', @@ -224,7 +247,7 @@ test('submit routes a unique complete Session name without asking', async () => assert.deepEqual(result, { kind: 'submitted', - strategyId: 'wh-r2.3-session-core-evidence', + strategyId: WORKHUB_ROUTING_STRATEGY_ID, requestId: 'request-exact', target: { sessionId: 'payment' }, turnId: 'turn-exact', @@ -333,7 +356,7 @@ test('submit asks the user when weak relevance matches more than one Session', a assert.deepEqual(result, { kind: 'clarification', - strategyId: 'wh-r2.3-session-core-evidence', + strategyId: WORKHUB_ROUTING_STRATEGY_ID, requestId: 'request-ambiguous', text: '继续处理重复问题', options: [ @@ -478,7 +501,7 @@ test('submit follows an unambiguous reference to the most recent Work', async () assert.deepEqual(result, { kind: 'submitted', - strategyId: 'wh-r2.3-session-core-evidence', + strategyId: WORKHUB_ROUTING_STRATEGY_ID, requestId: 'request-pronoun', target: { sessionId: 'payment' }, turnId: 'turn-2', @@ -487,6 +510,280 @@ test('submit follows an unambiguous reference to the most recent Work', async () assert.deepEqual(submitted, ['payment', 'payment']); }); +test('read seeds current and previous focus from pre-existing ordinary Sessions', async () => { + const submitted: string[] = []; + const sessions = port([ + session('login', { sessionName: '登录刷新令牌', updatedAt: 20 }), + session('payment', { sessionName: '支付回调幂等性', updatedAt: 30 }), + session('archived', { archived: true, updatedAt: 40 }), + ]); + sessions.submit = async (target) => { + submitted.push(target.sessionId); + return { turnId: `turn-${submitted.length}` }; + }; + const controller = createWorkHubController({ sessions }); + + await controller.read(); + const current = await controller.submit({ + requestId: 'request-current-seed', + text: '继续这个工作', + }); + const previous = await controller.submit({ + requestId: 'request-previous-seed', + text: '回到上一个工作', + }); + + assert.deepEqual(current.kind === 'submitted' ? current.target : undefined, { + sessionId: 'payment', + }); + assert.deepEqual(previous.kind === 'submitted' ? previous.target : undefined, { + sessionId: 'login', + }); + assert.deepEqual(submitted, ['payment', 'login']); +}); + +test('read prefers the Session active when WorkHub opens over raw recency', async () => { + const submitted: string[] = []; + const sessions = port([ + session('login', { sessionName: '登录刷新令牌', updatedAt: 20 }), + session('payment', { sessionName: '支付回调幂等性', updatedAt: 30 }), + ]); + sessions.submit = async (target) => { + submitted.push(target.sessionId); + return { turnId: 'turn-login' }; + }; + const controller = createWorkHubController({ sessions }); + + await controller.read({ focus: { sessionId: 'login' } }); + const result = await controller.submit({ + requestId: 'request-active-seed', + text: '继续这个工作', + }); + + assert.deepEqual(result.kind === 'submitted' ? result.target : undefined, { + sessionId: 'login', + }); + assert.deepEqual(submitted, ['login']); +}); + +test('a stale opening read cannot overwrite a newer WorkHub focus', async () => { + const pendingReads: Array<{ + resolve(value: WorkHubSessionFacts[]): void; + promise: Promise; + }> = []; + const sessions = port([]); + sessions.list = () => { + let resolve!: (value: WorkHubSessionFacts[]) => void; + const promise = new Promise((next) => { + resolve = next; + }); + pendingReads.push({ resolve, promise }); + return promise; + }; + const submitted: string[] = []; + sessions.submit = async (target) => { + submitted.push(target.sessionId); + return { turnId: 'turn-newer-focus' }; + }; + const controller = createWorkHubController({ sessions }); + const older = controller.read({ focus: { sessionId: 'payment' } }); + const newer = controller.read({ focus: { sessionId: 'login' } }); + const facts = [ + session('login', { updatedAt: 20 }), + session('payment', { updatedAt: 30 }), + ]; + + pendingReads[1]!.resolve(facts); + await newer; + pendingReads[0]!.resolve([facts[1]!]); + await older; + sessions.list = async () => facts; + const result = await controller.submit({ + requestId: 'request-after-stale-read', + text: '继续这个工作', + }); + + assert.deepEqual(result.kind === 'submitted' ? result.target : undefined, { + sessionId: 'login', + }); + assert.deepEqual(submitted, ['login']); +}); + +test('an unavailable opening focus falls back to recent routable Sessions', async () => { + const submitted: string[] = []; + const sessions = port([ + session('archived', { archived: true, updatedAt: 40 }), + session('login', { updatedAt: 20 }), + session('payment', { updatedAt: 30 }), + ]); + sessions.submit = async (target) => { + submitted.push(target.sessionId); + return { turnId: 'turn-fallback' }; + }; + const controller = createWorkHubController({ sessions }); + + await controller.read({ focus: { sessionId: 'archived' } }); + const result = await controller.submit({ + requestId: 'request-fallback-focus', + text: '继续这个工作', + }); + + assert.deepEqual(result.kind === 'submitted' ? result.target : undefined, { + sessionId: 'payment', + }); + assert.deepEqual(submitted, ['payment']); +}); + +test('focus falls back when the current Session is archived after WorkHub opens', async () => { + let catalog = [ + session('login', { updatedAt: 20 }), + session('payment', { updatedAt: 30 }), + ]; + const sessions = port(catalog); + sessions.list = async () => catalog; + const submitted: string[] = []; + sessions.submit = async (target) => { + submitted.push(target.sessionId); + return { turnId: 'turn-focus-fallback' }; + }; + const controller = createWorkHubController({ sessions }); + await controller.read(); + catalog = catalog.map((entry) => entry.target.sessionId === 'payment' + ? { ...entry, archived: true } + : entry); + + const result = await controller.submit({ + requestId: 'request-after-current-archive', + text: '继续这个工作', + }); + + assert.deepEqual(result.kind === 'submitted' ? result.target : undefined, { + sessionId: 'login', + }); + assert.deepEqual(submitted, ['login']); +}); + +test('resetVisitContext discards focus from a previous WorkHub mount', async () => { + let catalog = [ + session('login', { updatedAt: 20 }), + session('payment', { updatedAt: 30 }), + ]; + const sessions = port(catalog); + sessions.list = async () => catalog; + const submitted: string[] = []; + sessions.submit = async (target) => { + submitted.push(target.sessionId); + return { turnId: `turn-${submitted.length}` }; + }; + const controller = createWorkHubController({ sessions }); + await controller.read({ focus: { sessionId: 'login' } }); + controller.resetVisitContext(); + catalog = catalog.map((entry) => entry.target.sessionId === 'payment' + ? { ...entry, updatedAt: 40 } + : entry); + await controller.read(); + + const result = await controller.submit({ + requestId: 'request-after-remount', + text: '继续这个工作', + }); + + assert.deepEqual(result.kind === 'submitted' ? result.target : undefined, { + sessionId: 'payment', + }); +}); + +test('an in-flight submit cannot restore visit focus after WorkHub unmounts', async () => { + const sessions = port([ + session('login', { updatedAt: 20 }), + session('payment', { updatedAt: 30 }), + ]); + let signalSubmitStarted!: () => void; + const submitStarted = new Promise((resolve) => { + signalSubmitStarted = resolve; + }); + let finishSubmit!: (value: { turnId: string }) => void; + const pendingTurn = new Promise<{ turnId: string }>((resolve) => { + finishSubmit = resolve; + }); + sessions.submit = async () => { + signalSubmitStarted(); + return pendingTurn; + }; + const controller = createWorkHubController({ sessions }); + await controller.read({ focus: { sessionId: 'login' } }); + const inFlight = controller.submit({ + requestId: 'request-before-unmount', + text: '继续这个工作', + }); + await submitStarted; + controller.resetVisitContext(); + finishSubmit({ turnId: 'turn-login' }); + await inFlight; + + const submitted: string[] = []; + sessions.submit = async (target) => { + submitted.push(target.sessionId); + return { turnId: 'turn-after-remount' }; + }; + await controller.read(); + const result = await controller.submit({ + requestId: 'request-after-in-flight', + text: '继续这个工作', + }); + + assert.deepEqual(result.kind === 'submitted' ? result.target : undefined, { + sessionId: 'payment', + }); + assert.deepEqual(submitted, ['payment']); +}); + +test('an old submit resolves against the visit focus captured before an await', async () => { + const catalog = [ + session('login', { updatedAt: 20 }), + session('payment', { updatedAt: 30 }), + ]; + const sessions = port(catalog); + const controller = createWorkHubController({ sessions }); + await controller.read({ focus: { sessionId: 'login' } }); + + let signalListStarted!: () => void; + const listStarted = new Promise((resolve) => { + signalListStarted = resolve; + }); + let finishOldList!: (value: WorkHubSessionFacts[]) => void; + const oldList = new Promise((resolve) => { + finishOldList = resolve; + }); + let blockNextList = true; + sessions.list = async () => { + if (!blockNextList) return catalog; + blockNextList = false; + signalListStarted(); + return oldList; + }; + const submitted: string[] = []; + sessions.submit = async (target) => { + submitted.push(target.sessionId); + return { turnId: 'turn-' + target.sessionId }; + }; + + const oldSubmission = controller.submit({ + requestId: 'request-old-visit', + text: '继续这个工作', + }); + await listStarted; + controller.resetVisitContext(); + await controller.read({ focus: { sessionId: 'payment' } }); + finishOldList(catalog); + + const result = await oldSubmission; + assert.deepEqual(result.kind === 'submitted' ? result.target : undefined, { + sessionId: 'login', + }); + assert.deepEqual(submitted, ['login']); +}); + test('submit routes strong core evidence instead of reusing recent focus', async () => { const submitted: string[] = []; const sessions = port([ @@ -517,7 +814,7 @@ test('submit routes strong core evidence instead of reusing recent focus', async assert.deepEqual(result, { kind: 'submitted', - strategyId: 'wh-r2.3-session-core-evidence', + strategyId: WORKHUB_ROUTING_STRATEGY_ID, requestId: 'request-topic-shift', target: { sessionId: 'payment' }, turnId: 'turn-2', @@ -553,7 +850,7 @@ test('submit routes unique strong core evidence without asking', async () => { if (result.kind !== 'submitted') return; assert.deepEqual(result.target, { sessionId: 'login' }); assert.equal(result.evidence, 'core_entity'); - assert.equal(result.strategyId, 'wh-r2.3-session-core-evidence'); + assert.equal(result.strategyId, 'wh-r2.4-session-context-continuity'); assert.deepEqual(submitted, ['login']); }); @@ -579,14 +876,14 @@ test('submit ignores shared boilerplate when an executable request names a new t if (result.kind !== 'submitted') return; assert.deepEqual(result.target, { sessionId: 'payment-new' }); assert.equal(result.evidence, 'new_session'); - assert.equal(result.strategyId, 'wh-r2.3-session-core-evidence'); + assert.equal(result.strategyId, 'wh-r2.4-session-context-continuity'); assert.deepEqual(createdNames, ['检查支付回调重复投递']); }); -test('submit keeps a unique two-character clue behind clarification', async () => { +test('submit keeps a foreign two-character clue behind clarification', async () => { const sessions = port([ - session('login', { sessionName: '登录稳定性' }), - session('payment', { sessionName: '支付稳定性' }), + session('login', { sessionName: '登录稳定性', updatedAt: 10 }), + session('payment', { sessionName: '支付稳定性', updatedAt: 20 }), ]); const controller = createWorkHubController({ sessions }); @@ -596,7 +893,7 @@ test('submit keeps a unique two-character clue behind clarification', async () = }); assert.equal(result.kind, 'clarification'); - assert.equal(result.strategyId, 'wh-r2.3-session-core-evidence'); + assert.equal(result.strategyId, 'wh-r2.4-session-context-continuity'); }); test('submit treats explicit user uncertainty as clarification instead of a new Session', async () => { @@ -793,29 +1090,1177 @@ test('route correction never stops a root Turn that WorkHub only steered into', assert.deepEqual(stopped, []); }); -test('latest route correction wins for the same expression family', async () => { +test('first natural-language correction reroutes and stops the wrong WorkHub-owned Turn', async () => { + const submitted: string[] = []; + const stopped: Array<[string, string]> = []; const sessions = port([ - session('login', { sessionName: '登录稳定性' }), - session('payment', { sessionName: '支付稳定性' }), + session('login', { + sessionName: '登录稳定性', + latestResult: '刷新令牌过期导致重复登录', + updatedAt: 20, + }), + session('payment', { + sessionName: '支付稳定性', + latestResult: '支付回调重复投递', + updatedAt: 30, + }), ]); - sessions.submit = async (_target) => ({ turnId: 'turn' }); + 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: '继续这个工作,补充验收项', + }); + + const corrected = await controller.submit({ + requestId: 'request-natural-correction', + 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.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); + }; const controller = createWorkHubController({ sessions }); + await controller.read(); + await controller.submit({ + requestId: 'request-before-content-change', + text: '继续这个工作', + }); + + const result = await controller.submit({ + requestId: 'request-content-change', + text: '继续这个工作,Redis 配置不对,改成 Postgres', + }); + + assert.deepEqual(result.kind === 'submitted' ? result.target : undefined, { + sessionId: 'payment', + }); + 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: 'correction-login', - text: '继续白鹭点,列出验收项。', + requestId: 'request-owned-root', + text: '继续这个工作', + }); + await controller.submit({ + requestId: 'request-other-owned-root', + text: '先检查登录稳定性', explicitTarget: { sessionId: 'login' }, - correction: { from: { sessionId: 'payment' }, turnId: 'turn' }, }); await controller.submit({ - requestId: 'correction-payment', - text: '继续白鹭点,列出异常项。', + requestId: 'request-steer-owned-root', + text: '继续这个工作,补充测试点', explicitTarget: { sessionId: 'payment' }, - correction: { from: { sessionId: 'login' }, turnId: 'turn' }, }); - const result = await controller.submit({ - requestId: 'correction-latest', + 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 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: '继续白鹭点,补充回滚条件。', }); @@ -846,7 +2291,7 @@ test('waiting Session rejects a second root request without calling submit', asy assert.deepEqual(result, { kind: 'waiting', - strategyId: 'wh-r2.3-session-core-evidence', + strategyId: WORKHUB_ROUTING_STRATEGY_ID, requestId: 'request-waiting', text: '排查令牌过期重复登录问题:补充一条等待状态下的新请求。', target: { sessionId: 'login' }, @@ -935,7 +2380,7 @@ test('submit keeps unmatched non-executable conversation in WorkHub', async () = assert.deepEqual(result, { kind: 'discussion', - strategyId: 'wh-r2.3-session-core-evidence', + strategyId: WORKHUB_ROUTING_STRATEGY_ID, requestId: 'request-discussion', text: '你觉得统一入口最重要的价值是什么?', }); @@ -1000,7 +2445,7 @@ test('submit creates an ordinary Session for a clear unmatched executable goal', assert.deepEqual(result, { kind: 'submitted', - strategyId: 'wh-r2.3-session-core-evidence', + strategyId: WORKHUB_ROUTING_STRATEGY_ID, requestId: 'request-new-work', target: { sessionId: 'invoice-export' }, turnId: 'turn-invoice-export', 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 e96bebb817..719d9c81ca 100644 --- a/apps/desktop/src/main/__tests__/workhub-session-port.test.ts +++ b/apps/desktop/src/main/__tests__/workhub-session-port.test.ts @@ -27,6 +27,7 @@ import { projectWorkHubSessionTurns, type WorkHubDesktopSession, } from '../../renderer/workhub-session-port.js'; +import { WorkHubSessionSubmitError } from '../../renderer/workhub-controller.js'; function desktopSession( id: string, @@ -51,6 +52,48 @@ const unusedTranscripts = { }, }; +function transcriptsWith(messages: readonly StoredMessage[]) { + return { + open: async (sessionId: string, handler: (batch: DesktopTranscriptBatch) => void) => { + const parsed = JSON.parse(sessionId) as [string, string]; + const fragments = messages.map((message, identity) => { + const data = new TextEncoder().encode(JSON.stringify(message)); + return { + source: 'durable' as const, + identity, + order: null, + byteOffset: 0, + totalBytes: data.byteLength, + data, + }; + }); + handler({ + sessionId: parsed[1], + deliverySequence: 1, + generation: 'generation-reconcile', + hostEpoch: 'epoch-reconcile', + durableThrough: messages.length - 1, + fragments, + evictedDurableSequences: [], + completedOverlayMessageIds: [], + hasOlder: false, + hasNewer: false, + reset: true, + ready: true, + }); + return { + sessionId, + generation: 'generation-reconcile', + hostEpoch: 'epoch-reconcile', + readThroughMessageId: null, + loadBefore: async () => {}, + loadAround: async () => {}, + close: async () => {}, + }; + }, + }; +} + test('projects durable Session messages into an ordered WorkHub conversation', () => { const turns = projectWorkHubSessionTurns({ target: { sessionId: 'payment' }, @@ -310,6 +353,7 @@ test('desktop adapter projects Session catalog facts without owning copies', asy kind: 'ordinary', archived: false, state: 'running', + runningTurnIds: ['turn-running'], latestResult: '正在补充重复投递测试', updatedAt: 30, }, @@ -320,6 +364,7 @@ test('desktop adapter projects Session catalog facts without owning copies', asy kind: 'internal', archived: false, state: 'active', + runningTurnIds: [], updatedAt: 20, }, { @@ -329,6 +374,7 @@ test('desktop adapter projects Session catalog facts without owning copies', asy kind: 'ordinary', archived: false, state: 'waiting_for_user', + runningTurnIds: ['turn-waiting'], updatedAt: 15, }, { @@ -338,11 +384,41 @@ test('desktop adapter projects Session catalog facts without owning copies', asy kind: 'subagent', archived: false, state: 'active', + runningTurnIds: [], updatedAt: 10, }, ]); }); +test('desktop adapter preserves per-Host catalog coverage for ownership reconciliation', async () => { + const localSessionId = desktopSessionKey({ hostId: 'local-host', sessionId: 'local' }); + const adapter = createDesktopWorkHubSessionPort({ + transcripts: unusedTranscripts, + sessions: { + list: async () => [], + listWithCoverage: async () => ({ + sessions: [desktopSession(localSessionId)], + completeHostIds: ['local-host'], + }), + 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); +}); + test('desktop adapter delegates create, send, and invalidation to Session APIs', async () => { const calls: unknown[] = []; let onChanged: (() => void) | undefined; @@ -375,7 +451,8 @@ test('desktop adapter delegates create, send, and invalidation to Session APIs', }); const created = await adapter.create({ name: '实现导出发票 PDF 功能' }); - const turn = await adapter.submit(created.target, '实现导出发票 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(() => { @@ -417,11 +494,111 @@ test('desktop adapter preserves when Session delivery steered an existing root T }); assert.deepEqual( - await adapter.submit({ sessionId: 'busy' }, '补充已有执行流'), + 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' | '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'); + return { ok: false as const, reason: 'archived' as const }; + }, + stop: async () => {}, + subscribeChanges: () => () => {}, + }, + projectName: () => 'Maka', + newTurnId: () => 'reserved-turn', + }); + + 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({ 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 dc6479c082..27edc7eb97 100644 --- a/apps/desktop/src/main/__tests__/workhub-surface-flow.test.ts +++ b/apps/desktop/src/main/__tests__/workhub-surface-flow.test.ts @@ -31,6 +31,7 @@ import { import { boundedWorkHubTimelineText, createWorkHubController, + WORKHUB_ROUTING_STRATEGY_ID, type WorkHubController, type WorkHubSubmitInput, } from '../../renderer/workhub-controller.js'; @@ -79,14 +80,14 @@ test('surface keeps the Composer draft when routing fails or the target is waiti assert.equal(workHubSubmissionClearsDraft(undefined), false); assert.equal(workHubSubmissionClearsDraft({ kind: 'waiting', - strategyId: 'wh-r2.3-session-core-evidence', + strategyId: WORKHUB_ROUTING_STRATEGY_ID, requestId: 'waiting', text: '继续处理', target: { sessionId: 'payment' }, }), false); assert.equal(workHubSubmissionClearsDraft({ kind: 'discussion', - strategyId: 'wh-r2.3-session-core-evidence', + strategyId: WORKHUB_ROUTING_STRATEGY_ID, requestId: 'discussion', text: '先讨论方向', }), true); @@ -95,7 +96,7 @@ test('surface keeps the Composer draft when routing fails or the target is waiti test('surface disables correction after a request was steered into existing work', () => { const submission = { kind: 'submitted' as const, - strategyId: 'wh-r2.3-session-core-evidence' as const, + strategyId: WORKHUB_ROUTING_STRATEGY_ID, requestId: 'steered', target: { sessionId: 'payment' }, turnId: 'turn-existing', @@ -130,7 +131,7 @@ test('surface hides a rebuilt Session turn while the matching local turn is stil state: 'settled', outcome: { kind: 'submitted', - strategyId: 'wh-r2.3-session-core-evidence', + strategyId: WORKHUB_ROUTING_STRATEGY_ID, requestId: 'request-payment', target: { sessionId: 'payment' }, turnId: 'turn-payment', @@ -164,7 +165,7 @@ test('surface canonicalizes bounded text before suppressing a local duplicate', state: 'settled', outcome: { kind: 'submitted', - strategyId: 'wh-r2.3-session-core-evidence', + strategyId: WORKHUB_ROUTING_STRATEGY_ID, requestId: 'request-long', target: { sessionId: 'payment' }, turnId: 'turn-long', @@ -195,7 +196,7 @@ test('surface suppresses the newest matching projected steering turn', () => { state: 'settled', outcome: { kind: 'submitted', - strategyId: 'wh-r2.3-session-core-evidence', + strategyId: WORKHUB_ROUTING_STRATEGY_ID, requestId: 'request-new', target: { sessionId: 'payment' }, turnId: 'turn-payment', @@ -208,13 +209,14 @@ test('surface keeps clarification and successful routing in WorkHub', async () = const submissions: WorkHubSubmitInput[] = []; const controller: WorkHubController = { read: async () => ({ sessions: [], turns: [] }), + resetVisitContext: () => {}, subscribe: () => () => {}, submit: async (input) => { submissions.push(input); if (!input.explicitTarget) { return { kind: 'clarification', - strategyId: 'wh-r2.3-session-core-evidence', + strategyId: WORKHUB_ROUTING_STRATEGY_ID, requestId: input.requestId, text: input.text, options: [{ @@ -226,7 +228,7 @@ test('surface keeps clarification and successful routing in WorkHub', async () = } return { kind: 'submitted', - strategyId: 'wh-r2.3-session-core-evidence', + strategyId: WORKHUB_ROUTING_STRATEGY_ID, requestId: input.requestId, target: input.explicitTarget, turnId: 'turn-payment', @@ -256,10 +258,11 @@ test('surface keeps clarification and successful routing in WorkHub', async () = test('surface leaves discussion in WorkHub instead of creating a task view', async () => { const controller: WorkHubController = { read: async () => ({ sessions: [], turns: [] }), + resetVisitContext: () => {}, subscribe: () => () => {}, submit: async (input) => ({ kind: 'discussion', - strategyId: 'wh-r2.3-session-core-evidence', + strategyId: WORKHUB_ROUTING_STRATEGY_ID, requestId: input.requestId, text: input.text, }), diff --git a/apps/desktop/src/main/runtime-host-session-execution-ipc-main.ts b/apps/desktop/src/main/runtime-host-session-execution-ipc-main.ts index 3be6608f6a..bd961f7fdd 100644 --- a/apps/desktop/src/main/runtime-host-session-execution-ipc-main.ts +++ b/apps/desktop/src/main/runtime-host-session-execution-ipc-main.ts @@ -299,7 +299,9 @@ export function registerRuntimeHostSessionExecutionIpc( } const submitted = await deps.client.submitMessage({ sessionId, - messageId: newId(), + // Preserve the renderer's command identity in the durable message so + // a lost IPC reply can be reconciled as root-vs-steering later. + messageId: turnId, content: startInput.content, placement: "current_turn", }); diff --git a/apps/desktop/src/preload/bridge-contract.d.ts b/apps/desktop/src/preload/bridge-contract.d.ts index 226e27a324..d478543721 100644 --- a/apps/desktop/src/preload/bridge-contract.d.ts +++ b/apps/desktop/src/preload/bridge-contract.d.ts @@ -712,6 +712,10 @@ export interface MakaBridge { }; sessions: { list(filter?: SessionListFilter): Promise; + listWithCoverage(): Promise<{ + sessions: DesktopSessionSummary[]; + completeHostIds: string[]; + }>; create(input?: CreateSessionRequestInput): Promise; send( sessionId: string, diff --git a/apps/desktop/src/preload/preload.ts b/apps/desktop/src/preload/preload.ts index ab3725f04c..41839eac37 100644 --- a/apps/desktop/src/preload/preload.ts +++ b/apps/desktop/src/preload/preload.ts @@ -105,7 +105,10 @@ import type { import type { BotProvider } from '@maka/core/bot-chat-settings'; import type { BotOnboardingSnapshot, BotOnboardingStartInput } from '@maka/core/bot-onboarding'; import type { HealthSnapshot } from '@maka/core/health'; -import { collectRuntimeHostSessionCatalogs } from './runtime-host-session-catalog.js'; +import { + collectRuntimeHostSessionCatalogs, + collectRuntimeHostSessionCatalogsWithCoverage, +} from './runtime-host-session-catalog.js'; import type { ExecutionBoundaryReadModel, SandboxBoundaryResponse } from '@maka/core/sandbox-boundary'; import type { ActiveInteractionRequestEvent, @@ -819,6 +822,21 @@ async function listDesktopSessions( ); } +async function listDesktopSessionsWithCoverage(): Promise<{ + sessions: DesktopSessionSummary[]; + completeHostIds: string[]; +}> { + const scopes = await runtimeHostScopeList(); + return collectRuntimeHostSessionCatalogsWithCoverage( + scopes.map((scope) => ({ + hostId: scope.hostId, + sessions: ipcRenderer.invoke('sessions:list', scope) + .then((sessions: SessionCatalogSummary[]) => + sessions.map((session) => projectSessionSummary(scope, session))), + })), + ); +} + function sendActiveRuntimeHost(channel: string, ...args: unknown[]): void { void activeRuntimeHostRef() .then((scope) => ipcRenderer.send(channel, scope, ...args)) @@ -1492,6 +1510,9 @@ const makaBridge = { list(filter?: SessionListFilter): Promise { return listDesktopSessions(filter); }, + listWithCoverage() { + return listDesktopSessionsWithCoverage(); + }, /** * The single session-creation channel (#1433). `mode` names a * product intent — main derives the permission boundary, name and diff --git a/apps/desktop/src/preload/runtime-host-session-catalog.ts b/apps/desktop/src/preload/runtime-host-session-catalog.ts index a23f9a6c37..91b721d6b6 100644 --- a/apps/desktop/src/preload/runtime-host-session-catalog.ts +++ b/apps/desktop/src/preload/runtime-host-session-catalog.ts @@ -19,6 +19,35 @@ import type { DesktopSessionSummary } from './bridge-contract.js'; +export interface RuntimeHostSessionCatalogRequest { + readonly hostId: string; + readonly sessions: Promise; +} + +export interface RuntimeHostSessionCatalogCoverage { + readonly sessions: DesktopSessionSummary[]; + readonly completeHostIds: string[]; +} + +export async function collectRuntimeHostSessionCatalogsWithCoverage( + requests: readonly RuntimeHostSessionCatalogRequest[], +): Promise { + const results = await Promise.allSettled(requests.map((request) => request.sessions)); + const fulfilled = results.flatMap((result, index) => result.status === 'fulfilled' + ? [{ hostId: requests[index]!.hostId, sessions: result.value }] + : []); + if (requests.length > 0 && fulfilled.length === 0) { + throw new AggregateError( + results.flatMap((result) => result.status === 'rejected' ? [result.reason] : []), + 'Every Runtime Host Session Catalog request failed', + ); + } + return { + sessions: sortSessionCatalogs(fulfilled.flatMap((entry) => entry.sessions)), + completeHostIds: fulfilled.map((entry) => entry.hostId), + }; +} + export async function collectRuntimeHostSessionCatalogs( requests: readonly Promise[], ): Promise { @@ -30,7 +59,11 @@ export async function collectRuntimeHostSessionCatalogs( 'Every Runtime Host Session Catalog request failed', ); } - return groups.flat().sort((left, right) => { + return sortSessionCatalogs(groups.flat()); +} + +function sortSessionCatalogs(sessions: DesktopSessionSummary[]): DesktopSessionSummary[] { + return sessions.sort((left, right) => { if (left.activityAt === undefined || right.activityAt === undefined) { throw new Error('Runtime Host Session Catalog activity is unavailable'); } diff --git a/apps/desktop/src/renderer/app-shell.tsx b/apps/desktop/src/renderer/app-shell.tsx index 30ed340d22..c599125ae6 100644 --- a/apps/desktop/src/renderer/app-shell.tsx +++ b/apps/desktop/src/renderer/app-shell.tsx @@ -1476,14 +1476,21 @@ function AppShellContent({ captureActiveComposerClaim, }); refreshProjectSkillsRef.current = moduleHub.commands.refreshProjectSkills; - const workHubController = useMemo(() => createWorkHubController({ - sessions: createDesktopWorkHubSessionPort({ - sessions: window.maka.sessions, - transcripts: window.maka.transcripts, - projectName: (projectId) => projects.find((project) => project.id === projectId)?.name, - newTurnId: () => crypto.randomUUID(), - }), - }), [projects]); + const workHubProjectsRef = useRef(projects); + workHubProjectsRef.current = projects; + const workHubControllerRef = useRef | null>(null); + if (!workHubControllerRef.current) { + workHubControllerRef.current = createWorkHubController({ + sessions: createDesktopWorkHubSessionPort({ + sessions: window.maka.sessions, + transcripts: window.maka.transcripts, + projectName: (projectId) => + workHubProjectsRef.current.find((project) => project.id === projectId)?.name, + newTurnId: () => crypto.randomUUID(), + }), + }); + } + const workHubController = workHubControllerRef.current; // Where a NEW chat starts. Built unconditionally and handed to the composer, // which renders it only while no session owns it — the project is fixed once // the first message creates one, so there is nothing to pick after that. @@ -2802,6 +2809,7 @@ function AppShellContent({ ) : ( diff --git a/apps/desktop/src/renderer/styles/workhub.css b/apps/desktop/src/renderer/styles/workhub.css index 369a71c038..ff666d9d98 100644 --- a/apps/desktop/src/renderer/styles/workhub.css +++ b/apps/desktop/src/renderer/styles/workhub.css @@ -91,6 +91,15 @@ flex-direction: column; } +.workhub-loading-user { + display: flex; + justify-content: flex-end; +} + +.workhub-loading-turn { + min-height: 172px; +} + .workhub-turn { display: flex; flex-direction: column; diff --git a/apps/desktop/src/renderer/workhub-controller.ts b/apps/desktop/src/renderer/workhub-controller.ts index 514d47794d..e382ee7fd5 100644 --- a/apps/desktop/src/renderer/workhub-controller.ts +++ b/apps/desktop/src/renderer/workhub-controller.ts @@ -47,11 +47,13 @@ export interface WorkHubSessionFacts { kind: 'ordinary' | 'internal' | 'subagent'; archived: boolean; state: WorkHubSessionState; + /** Authoritative live Turn IDs when the Session catalog provides them. */ + runningTurnIds?: readonly string[]; latestResult?: string; updatedAt: number; } -export type WorkHubSessionSummary = Omit; +export type WorkHubSessionSummary = Omit; export type WorkHubProjectedTurnState = 'running' | 'completed' | 'aborted' | 'failed'; @@ -84,14 +86,20 @@ export interface WorkHubSubmitInput { requestId: string; text: string; explicitTarget?: WorkHubSessionTarget; - correction?: { - from: WorkHubSessionTarget; - turnId: string; - steered?: true; - }; + correction?: WorkHubCorrectionContext; +} + +export interface WorkHubCorrectionContext { + from: WorkHubSessionTarget; + turnId?: string; + steered?: true; +} + +export interface WorkHubReadInput { + focus?: WorkHubSessionTarget; } -export const WORKHUB_ROUTING_STRATEGY_ID = 'wh-r2.3-session-core-evidence' as const; +export const WORKHUB_ROUTING_STRATEGY_ID = 'wh-r2.4-session-context-continuity' as const; export type WorkHubRoutingStrategyId = typeof WORKHUB_ROUTING_STRATEGY_ID; export type WorkHubSubmission = ( @@ -109,6 +117,7 @@ export type WorkHubSubmission = ( requestId: string; text: string; options: Array>; + correction?: WorkHubCorrectionContext; } | { kind: 'discussion'; @@ -129,6 +138,14 @@ export type WorkHubSubmission = ( */ export interface WorkHubSessionPort { list(): Promise; + /** + * Lists Sessions with per-target catalog coverage. A target missing from a + * partial multi-Host list is not authoritatively absent. + */ + listCatalog?(): Promise<{ + sessions: WorkHubSessionFacts[]; + isCompleteFor(target: WorkHubSessionTarget): boolean; + }>; /** * Rebuilds a bounded recent conversation from the authoritative Session * transcripts. Missing transcripts are omitted rather than copied elsewhere. @@ -142,41 +159,448 @@ export interface WorkHubSessionPort { targets: readonly WorkHubSessionTarget[], ): Promise>; create(input: { name: string }): Promise; + reserveTurnId(): string; submit( target: WorkHubSessionTarget, text: string, + turnId: string, ): Promise<{ turnId: string; steered?: true }>; + reconcileSubmission( + target: WorkHubSessionTarget, + reservedTurnId: string, + ): Promise< + | { kind: 'root'; turnId: string } + | { kind: 'steered' } + | { kind: 'unknown' } + >; stop(target: WorkHubSessionTarget, expectedTurnId: string): Promise; subscribe(handler: () => void): () => void; } +export class WorkHubSessionSubmitError extends Error { + constructor( + message: string, + readonly admission: 'rejected' | 'unknown', + options?: ErrorOptions, + ) { + super(message, options); + this.name = 'WorkHubSessionSubmitError'; + } +} + export interface WorkHubController { - read(): Promise; + read(input?: WorkHubReadInput): Promise; submit(input: WorkHubSubmitInput): Promise; subscribe(handler: () => void): () => void; + resetVisitContext(): void; +} + +const MAX_TRACKED_WORKHUB_ROOTS = 32; + +interface WorkHubRootOwnership { + order: number; + turnId: string; +} + +interface WorkHubPendingAdmission extends WorkHubRootOwnership { + state: 'in_flight' | 'uncertain'; +} + +interface WorkHubOwnershipTombstone { + order: number; + stoppedTurnIds: Set; } export function createWorkHubController(deps: { sessions: WorkHubSessionPort; }): WorkHubController { - const routePolicy = createWorkHubRoutePolicy(); + let routePolicy = createWorkHubRoutePolicy(); + let focusReadVersion = 0; + let pendingFocusReadVersion: number | undefined; + const confirmedOwnershipBySessionId = new Map(); + const pendingAdmissionsBySessionId = new Map(); + const ownershipTombstoneBySessionId = new Map(); + const stopAttemptByTurn = new Map>(); + const stopOperationCountBySessionId = new Map(); + let ownershipRevision = 0; + const reconcileFocus = ( + policy: ReturnType, + sessions: readonly WorkHubSessionFacts[], + ) => { + policy.initializeFocus(sessions + .filter((session) => session.kind === 'ordinary' && !session.archived) + .sort((left, right) => right.updatedAt - left.updatedAt) + .map((session) => session.target)); + }; + const correctionFor = (from: WorkHubSessionTarget): WorkHubCorrectionContext => { + const confirmed = confirmedOwnershipBySessionId.get(from.sessionId); + const pending = pendingAdmissionsBySessionId.get(from.sessionId); + const turnId = confirmed?.turnId ?? pending?.at(-1)?.turnId; + if (!turnId) return { from }; + return { + from, + turnId, + }; + }; + const pendingAdmissions = (sessionId: string): WorkHubPendingAdmission[] => + pendingAdmissionsBySessionId.get(sessionId) ?? []; + const setPendingAdmissions = ( + sessionId: string, + pending: WorkHubPendingAdmission[], + ) => { + ownershipRevision += 1; + if (pending.length === 0) { + pendingAdmissionsBySessionId.delete(sessionId); + return; + } + pendingAdmissionsBySessionId.set( + sessionId, + [...pending].sort((left, right) => left.order - right.order), + ); + }; + const trackedRootCount = () => { + let pendingCount = 0; + for (const pending of pendingAdmissionsBySessionId.values()) { + pendingCount += pending.length; + } + return confirmedOwnershipBySessionId.size + pendingCount; + }; + const maybeRetireTombstone = (sessionId: string) => { + const tombstone = ownershipTombstoneBySessionId.get(sessionId); + if (!tombstone) return; + if ((stopOperationCountBySessionId.get(sessionId) ?? 0) > 0) return; + if (pendingAdmissions(sessionId).some((candidate) => candidate.order <= tombstone.order)) { + return; + } + ownershipTombstoneBySessionId.delete(sessionId); + }; + const readCatalog = async () => { + const revisionAtStart = ownershipRevision; + const catalog = deps.sessions.listCatalog + ? await deps.sessions.listCatalog() + : { + sessions: await deps.sessions.list(), + isCompleteFor: () => false, + }; + return { + catalog, + // A catalog request that overlapped an ownership mutation may describe + // the state before that mutation. It remains useful for projection, but + // it must not authoritatively prune newer ownership or admissions. + allowAuthoritativePruning: revisionAtStart === ownershipRevision, + }; + }; + const reconcileConfirmedOwnership = (catalog: { + sessions: readonly WorkHubSessionFacts[]; + isCompleteFor(target: WorkHubSessionTarget): boolean; + }, allowAuthoritativePruning: boolean) => { + if (!allowAuthoritativePruning) return; + const { sessions } = catalog; + const sessionById = new Map(sessions.map((session) => [session.target.sessionId, session])); + for (const [sessionId, ownership] of confirmedOwnershipBySessionId) { + const session = sessionById.get(sessionId); + if ( + (!session && catalog.isCompleteFor({ sessionId })) || + session?.archived || + (session?.runningTurnIds !== undefined && + !session.runningTurnIds.includes(ownership.turnId)) + ) { + if (confirmedOwnershipBySessionId.delete(sessionId)) { + ownershipRevision += 1; + } + } + } + }; + const storeOwnershipTombstone = ( + sessionId: string, + order: number, + stoppedTurnIds: Iterable = [], + ) => { + const existing = ownershipTombstoneBySessionId.get(sessionId); + if (existing && existing.order > order) return; + const stopped = new Set(existing?.order === order ? existing.stoppedTurnIds : []); + for (const turnId of stoppedTurnIds) stopped.add(turnId); + ownershipTombstoneBySessionId.set(sessionId, { + order, + stoppedTurnIds: stopped, + }); + }; + const reserveOwnedRoot = ( + target: WorkHubSessionTarget, + turnId: string, + order: number, + ) => { + if (trackedRootCount() >= MAX_TRACKED_WORKHUB_ROOTS) { + throw new Error('WorkHub has too many unresolved root submissions'); + } + setPendingAdmissions(target.sessionId, [ + ...pendingAdmissions(target.sessionId), + { order, turnId, state: 'in_flight' }, + ]); + }; + const removePendingRoot = ( + target: WorkHubSessionTarget, + reservedTurnId: string, + order: number, + ) => { + setPendingAdmissions( + target.sessionId, + pendingAdmissions(target.sessionId).filter((candidate) => + candidate.order !== order || candidate.turnId !== reservedTurnId), + ); + }; + const markPendingRootUncertain = ( + target: WorkHubSessionTarget, + reservedTurnId: string, + order: number, + ) => { + setPendingAdmissions( + target.sessionId, + pendingAdmissions(target.sessionId).map((candidate) => + candidate.order === order && candidate.turnId === reservedTurnId + ? { ...candidate, state: 'uncertain' } + : candidate), + ); + }; + const attemptStop = ( + target: WorkHubSessionTarget, + turnId: string, + ): Promise => { + const key = `${target.sessionId}\0${turnId}`; + const existing = stopAttemptByTurn.get(key); + if (existing) return existing; + stopOperationCountBySessionId.set( + target.sessionId, + (stopOperationCountBySessionId.get(target.sessionId) ?? 0) + 1, + ); + const stopping = deps.sessions.stop(target, turnId).finally(() => { + stopAttemptByTurn.delete(key); + const remaining = (stopOperationCountBySessionId.get(target.sessionId) ?? 1) - 1; + if (remaining === 0) { + stopOperationCountBySessionId.delete(target.sessionId); + } else { + stopOperationCountBySessionId.set(target.sessionId, remaining); + } + }); + stopAttemptByTurn.set(key, stopping); + return stopping; + }; + const settleOwnedRoot = async ( + target: WorkHubSessionTarget, + reservedTurnId: string, + turn: { turnId: string; steered?: true }, + order: number, + ) => { + const tombstone = ownershipTombstoneBySessionId.get(target.sessionId); + let stopped = false; + let stopFailure: unknown; + if (!turn.steered && tombstone && tombstone.order >= order) { + stopped = tombstone.stoppedTurnIds.has(turn.turnId); + if (!stopped) { + try { + const priorStopAttempt = stopAttemptByTurn.get( + `${target.sessionId}\0${turn.turnId}`, + ); + if (priorStopAttempt) { + try { + await priorStopAttempt; + } catch { + // Admission is new evidence. Retry against the admitted root even + // when the earlier pre-admission Stop failed or observed nothing. + } + } + await attemptStop(target, turn.turnId); + const currentBarrier = ownershipTombstoneBySessionId.get(target.sessionId); + if (currentBarrier && currentBarrier.order >= order) { + storeOwnershipTombstone(target.sessionId, currentBarrier.order, [turn.turnId]); + } + stopped = true; + } catch (error) { + stopFailure = error; + } + } + } + removePendingRoot(target, reservedTurnId, order); + if (!turn.steered && !stopped) { + const confirmed = confirmedOwnershipBySessionId.get(target.sessionId); + if (!confirmed || confirmed.order <= order) { + confirmedOwnershipBySessionId.set(target.sessionId, { + order, + turnId: turn.turnId, + }); + ownershipRevision += 1; + } + } + maybeRetireTombstone(target.sessionId); + if (stopFailure) throw stopFailure; + }; + const releasePendingRoot = ( + target: WorkHubSessionTarget, + reservedTurnId: string, + order: number, + ) => { + removePendingRoot(target, reservedTurnId, order); + maybeRetireTombstone(target.sessionId); + }; + const reconcilePendingRoot = async ( + target: WorkHubSessionTarget, + reservedTurnId: string, + order: number, + ): Promise => { + const reconciliation = await deps.sessions.reconcileSubmission(target, reservedTurnId); + if (reconciliation.kind === 'unknown') return false; + await settleOwnedRoot( + target, + reservedTurnId, + reconciliation.kind === 'steered' + ? { turnId: reservedTurnId, steered: true } + : { turnId: reconciliation.turnId }, + order, + ); + return true; + }; + const reconcileUncertainAdmissions = (catalog: { + sessions: readonly WorkHubSessionFacts[]; + isCompleteFor(target: WorkHubSessionTarget): boolean; + }, allowAuthoritativePruning: boolean): Promise | undefined => { + const sessionById = new Map( + catalog.sessions.map((session) => [session.target.sessionId, session]), + ); + const uncertain = [...pendingAdmissionsBySessionId.entries()] + .flatMap(([sessionId, pending]) => pending + .filter((candidate) => candidate.state === 'uncertain') + .map((candidate) => ({ + target: { sessionId }, + ...candidate, + }))); + if (uncertain.length === 0) return undefined; + return Promise.all(uncertain.map(async ({ target, turnId, order }) => { + const session = sessionById.get(target.sessionId); + if ( + allowAuthoritativePruning && + (session?.archived || (!session && catalog.isCompleteFor(target))) + ) { + releasePendingRoot(target, turnId, order); + return; + } + try { + await reconcilePendingRoot(target, turnId, order); + } catch { + // Failed reconciliation preserves the pending or confirmed ownership. + } + })).then(() => undefined); + }; + const assertSubmissionBarrierOpen = (target: WorkHubSessionTarget) => { + maybeRetireTombstone(target.sessionId); + const tombstone = ownershipTombstoneBySessionId.get(target.sessionId); + const stopCount = stopOperationCountBySessionId.get(target.sessionId) ?? 0; + const pendingBarrier = tombstone && pendingAdmissions(target.sessionId) + .some((candidate) => candidate.order <= tombstone.order); + if (stopCount > 0 || pendingBarrier) { + throw new Error('WorkHub is still reconciling a correction for this Session'); + } + }; + const stopOwnedRoots = async ( + correction: WorkHubCorrectionContext, + order: number, + ) => { + if (correction.steered) return; + const confirmed = confirmedOwnershipBySessionId.get(correction.from.sessionId); + const pending = pendingAdmissions(correction.from.sessionId); + const turnIds = new Set(); + const unconfirmedTurnIds = new Set(); + if (correction.turnId) turnIds.add(correction.turnId); + if (confirmed && confirmed.order < order) { + turnIds.add(confirmed.turnId); + } + for (const candidate of pending) { + if (candidate.order < order) { + turnIds.add(candidate.turnId); + unconfirmedTurnIds.add(candidate.turnId); + } + } + if (turnIds.size === 0) return; + // Publish only the order barrier before awaiting Host acknowledgements. + // Individual IDs become tombstoned only after their Stop succeeds. + storeOwnershipTombstone(correction.from.sessionId, order); + const failures: unknown[] = []; + await Promise.all([...turnIds].map(async (turnId) => { + try { + await attemptStop(correction.from, turnId); + const barrier = ownershipTombstoneBySessionId.get(correction.from.sessionId); + if (barrier && barrier.order >= order && !unconfirmedTurnIds.has(turnId)) { + storeOwnershipTombstone(correction.from.sessionId, barrier.order, [turnId]); + } + const owned = confirmedOwnershipBySessionId.get(correction.from.sessionId); + if (owned && owned.order < order && owned.turnId === turnId) { + confirmedOwnershipBySessionId.delete(correction.from.sessionId); + ownershipRevision += 1; + } + } catch (error) { + failures.push(error); + } + })); + maybeRetireTombstone(correction.from.sessionId); + if (failures.length > 0) throw failures[0]; + }; return { subscribe(handler) { return deps.sessions.subscribe(handler); }, - async read() { - const facts = await deps.sessions.list(); - const ordinary = facts - .filter((session) => session.kind === 'ordinary') - .sort((left, right) => right.updatedAt - left.updatedAt); - return { - sessions: ordinary - .map(({ kind: _kind, ...session }) => session), - turns: await deps.sessions.recentTurns(ordinary.map((session) => session.target)), - }; + async read(input) { + const readPolicy = routePolicy; + let readFocusVersion = focusReadVersion; + if (input?.focus) { + readFocusVersion = ++focusReadVersion; + pendingFocusReadVersion = readFocusVersion; + readPolicy.rememberTarget(input.focus); + } + try { + const { catalog, allowAuthoritativePruning } = + await readCatalog(); + reconcileConfirmedOwnership(catalog, allowAuthoritativePruning); + const reconciliation = reconcileUncertainAdmissions( + catalog, + allowAuthoritativePruning, + ); + if (reconciliation) await reconciliation; + const facts = catalog.sessions; + const ordinary = facts + .filter((session) => session.kind === 'ordinary') + .sort((left, right) => right.updatedAt - left.updatedAt); + if ( + readFocusVersion === focusReadVersion && + (input?.focus || pendingFocusReadVersion === undefined) + ) { + reconcileFocus(readPolicy, facts); + } + return { + sessions: ordinary + .map(({ kind: _kind, runningTurnIds: _runningTurnIds, ...session }) => session), + turns: await deps.sessions.recentTurns(ordinary.map((session) => session.target)), + }; + } finally { + if (input?.focus && pendingFocusReadVersion === readFocusVersion) { + pendingFocusReadVersion = undefined; + } + } }, async submit(input) { - const sessions = await deps.sessions.list(); + const submissionPolicy = routePolicy; + // Reserve the order synchronously, before any await. Corrections are + // learned only after successful delivery, but their precedence follows + // user submission order rather than network completion order. + const submissionOrder = submissionPolicy.reserveSubmissionOrder(); + const { catalog, allowAuthoritativePruning } = + await readCatalog(); + reconcileConfirmedOwnership(catalog, allowAuthoritativePruning); + const reconciliation = reconcileUncertainAdmissions( + catalog, + allowAuthoritativePruning, + ); + if (reconciliation) await reconciliation; + const sessions = catalog.sessions; + reconcileFocus(submissionPolicy, sessions); const ordinary = sessions.filter((session) => session.kind === 'ordinary'); // Archived Sessions remain visible as historical work, but Runtime Host // rejects new root Turns for them. Never offer one as a routing target. @@ -184,7 +608,7 @@ export function createWorkHubController(deps: { const routingEvidence = input.explicitTarget ? [] : await deps.sessions.routingEvidence(routable.map((session) => session.target)); - const decision = routePolicy.resolve({ + const decision = submissionPolicy.resolve({ text: input.text, sessions: routable, originPromptBySessionId: new Map( @@ -193,6 +617,9 @@ export function createWorkHubController(deps: { ...(input.explicitTarget ? { explicitTarget: input.explicitTarget } : {}), }); if (decision.kind === 'clarification') { + const correction = decision.correctedFrom + ? correctionFor(decision.correctedFrom) + : undefined; return { kind: 'clarification', strategyId: WORKHUB_ROUTING_STRATEGY_ID, @@ -203,6 +630,7 @@ export function createWorkHubController(deps: { projectName: session.projectName, sessionName: session.sessionName, })), + ...(correction ? { correction } : {}), }; } if (decision.kind === 'discussion') { @@ -215,6 +643,9 @@ export function createWorkHubController(deps: { } let target: WorkHubSessionTarget; let evidence: Extract['evidence']; + const correction = input.correction ?? (decision.kind === 'target' && decision.correctedFrom + ? correctionFor(decision.correctedFrom) + : undefined); if (decision.kind === 'new_session') { const created = await deps.sessions.create({ name: workHubNewSessionName(input.text) }); if (created.kind !== 'ordinary') { @@ -224,7 +655,7 @@ export function createWorkHubController(deps: { evidence = 'new_session'; } else { target = decision.target; - evidence = input.correction ? 'route_correction' : decision.evidence; + evidence = correction ? 'route_correction' : decision.evidence; } const targetSession = routable.find( (session) => session.target.sessionId === target.sessionId, @@ -241,12 +672,37 @@ export function createWorkHubController(deps: { target, }; } - if (input.correction && !input.correction.steered) { - await deps.sessions.stop(input.correction.from, input.correction.turnId); + if (correction) { + await stopOwnedRoots(correction, submissionOrder); + } + assertSubmissionBarrierOpen(target); + const reservedTurnId = deps.sessions.reserveTurnId(); + reserveOwnedRoot(target, reservedTurnId, submissionOrder); + let turn: { turnId: string; steered?: true }; + try { + turn = await deps.sessions.submit(target, input.text, reservedTurnId); + } catch (error) { + if ( + error instanceof WorkHubSessionSubmitError && + error.admission === 'rejected' + ) { + releasePendingRoot(target, reservedTurnId, submissionOrder); + } else { + markPendingRootUncertain(target, reservedTurnId, submissionOrder); + try { + await reconcilePendingRoot(target, reservedTurnId, submissionOrder); + } catch { + // The original delivery error remains primary. Reconciliation keeps + // any unresolved admission reachable for a later read/correction. + } + } + throw error; + } + await settleOwnedRoot(target, reservedTurnId, turn, submissionOrder); + submissionPolicy.rememberTarget(target); + if (correction) { + submissionPolicy.rememberCorrection(input.text, target, submissionOrder); } - const turn = await deps.sessions.submit(target, input.text); - routePolicy.rememberTarget(target); - if (input.correction) routePolicy.rememberCorrection(input.text, target); return { kind: 'submitted', strategyId: WORKHUB_ROUTING_STRATEGY_ID, @@ -255,8 +711,13 @@ export function createWorkHubController(deps: { turnId: turn.turnId, ...(turn.steered ? { steered: true as const } : {}), evidence, - ...(input.correction ? { correctedFrom: input.correction.from } : {}), + ...(correction ? { correctedFrom: correction.from } : {}), }; }, + resetVisitContext() { + focusReadVersion += 1; + pendingFocusReadVersion = undefined; + routePolicy = routePolicy.newVisit(); + }, }; } diff --git a/apps/desktop/src/renderer/workhub-route-policy.ts b/apps/desktop/src/renderer/workhub-route-policy.ts index 570031ba16..7edb39b5ad 100644 --- a/apps/desktop/src/renderer/workhub-route-policy.ts +++ b/apps/desktop/src/renderer/workhub-route-policy.ts @@ -34,10 +34,12 @@ export type WorkHubRouteDecision = kind: 'target'; target: WorkHubSessionTarget; evidence: WorkHubRouteEvidence; + correctedFrom?: WorkHubSessionTarget; } | { kind: 'clarification'; options: WorkHubSessionFacts[]; + correctedFrom?: WorkHubSessionTarget; } | { kind: 'discussion' } | { kind: 'new_session' }; @@ -49,8 +51,11 @@ export interface WorkHubRoutePolicy { originPromptBySessionId: ReadonlyMap; explicitTarget?: WorkHubSessionTarget; }): WorkHubRouteDecision; + initializeFocus(targets: readonly WorkHubSessionTarget[]): void; + newVisit(): WorkHubRoutePolicy; rememberTarget(target: WorkHubSessionTarget): void; - rememberCorrection(text: string, target: WorkHubSessionTarget): void; + reserveSubmissionOrder(): number; + rememberCorrection(text: string, target: WorkHubSessionTarget, order: number): void; } export function workHubNewSessionName(text: string): string { @@ -76,6 +81,11 @@ interface RouteCorrection { sequence: number; } +interface RouteCorrectionMemory { + corrections: RouteCorrection[]; + sequence: number; +} + const MAX_ROUTE_CORRECTIONS = 32; const MIN_EXACT_SESSION_NAME_LENGTH = 2; const MIN_CORRECTION_TERM_LENGTH = 3; @@ -88,16 +98,21 @@ const MAX_UNCERTAINTY_OPTIONS = 5; const MAX_RELATED_CLARIFICATION_OPTIONS = 4; /** - * Deep routing module for R2.3. + * Deep routing module for R2.4. * * It owns only transient inference context. Session identity, transcript, * execution state, and recovery continue to come from the Session port. */ export function createWorkHubRoutePolicy(): WorkHubRoutePolicy { + return createWorkHubRoutePolicyVisit({ corrections: [], sequence: 0 }); +} + +function createWorkHubRoutePolicyVisit( + correctionMemory: RouteCorrectionMemory, +): WorkHubRoutePolicy { let currentFocus: WorkHubSessionTarget | undefined; let previousFocus: WorkHubSessionTarget | undefined; - const corrections: RouteCorrection[] = []; - let correctionSequence = 0; + const corrections = correctionMemory.corrections; return { resolve({ text, sessions, originPromptBySessionId, explicitTarget }) { @@ -109,17 +124,58 @@ export function createWorkHubRoutePolicy(): WorkHubRoutePolicy { return { kind: 'new_session' }; } - const exact = sessions.map((session) => { - const qualifiedName = `${session.projectName}/${session.sessionName}`; + const correctionText = naturalCorrectionTargetText(text); + const correctedFrom = currentFocus && sessions.some((session) => + session.target.sessionId === currentFocus?.sessionId) + ? currentFocus + : undefined; + if ( + !correctionText && + correctedFrom && + looksLikeRecentFocus(text) && + looksLikeContentReplacement(text) + ) { + return { kind: 'target', target: correctedFrom, evidence: 'recent_focus' }; + } + if (correctionText && correctedFrom) { + const alternatives = sessions.filter((session) => + session.target.sessionId !== correctedFrom.sessionId); + const exactCorrection = rankExactSessions(correctionText, alternatives); + if ( + exactCorrection[0] && + exactCorrection[0].matchLength > (exactCorrection[1]?.matchLength ?? 0) + ) { + return { + kind: 'target', + target: exactCorrection[0].session.target, + evidence: 'route_correction', + correctedFrom, + }; + } + const relatedCorrection = rankRelatedSessions( + correctionText, + alternatives, + originPromptBySessionId, + ); + if (relatedCorrection.length === 1) { + return { + kind: 'target', + target: relatedCorrection[0]!.session.target, + evidence: 'route_correction', + correctedFrom, + }; + } + const options = relatedCorrection.length > 1 + ? relatedCorrection.map(({ session }) => session) + : alternatives.sort((left, right) => right.updatedAt - left.updatedAt); return { - session, - matchLength: Math.max( - exactIdentityMatchLength(text, qualifiedName), - exactIdentityMatchLength(text, session.sessionName), - ), + kind: 'clarification', + options: options.slice(0, MAX_UNCERTAINTY_OPTIONS), + correctedFrom, }; - }).filter(({ matchLength }) => matchLength >= MIN_EXACT_SESSION_NAME_LENGTH) - .sort((left, right) => right.matchLength - left.matchLength); + } + + const exact = rankExactSessions(text, sessions); if (exact[0] && exact[0].matchLength > (exact[1]?.matchLength ?? 0)) { return { kind: 'target', @@ -152,19 +208,28 @@ export function createWorkHubRoutePolicy(): WorkHubRoutePolicy { const previousReference = looksLikePreviousFocus(text); const currentReference = !previousReference && looksLikeRecentFocus(text); - const focused = previousReference + const focusCandidate = previousReference ? previousFocus : currentReference ? currentFocus : undefined; + const focused = focusCandidate && sessions.some((session) => + session.target.sessionId === focusCandidate.sessionId) + ? focusCandidate + : undefined; const strongEvidenceElsewhere = focused ? related.some(({ session, strongEvidence }) => session.target.sessionId !== focused.sessionId && strongEvidence) : false; + const weakEvidenceElsewhere = focused + ? related.some(({ session }) => session.target.sessionId !== focused.sessionId) + : false; const ambiguousEvidence = related.length > 1; if ( focused && - (previousReference || (!strongEvidenceElsewhere && !ambiguousEvidence)) + (previousReference || ( + !strongEvidenceElsewhere && !weakEvidenceElsewhere && !ambiguousEvidence + )) ) { return { kind: 'target', target: focused, evidence: 'recent_focus' }; } @@ -194,19 +259,63 @@ export function createWorkHubRoutePolicy(): WorkHubRoutePolicy { } return looksExecutable(text) ? { kind: 'new_session' } : { kind: 'discussion' }; }, + initializeFocus(targets) { + const ordered = targets.filter((target, index) => + targets.findIndex((candidate) => candidate.sessionId === target.sessionId) === index); + const first = ordered[0]; + if (!first) return; + const available = new Set(ordered.map((target) => target.sessionId)); + if (currentFocus && !available.has(currentFocus.sessionId)) { + currentFocus = first; + previousFocus = ordered[1]; + return; + } + if (!currentFocus) { + currentFocus = first; + previousFocus = ordered[1]; + return; + } + if (!previousFocus || !available.has(previousFocus.sessionId)) { + previousFocus = ordered.find((target) => target.sessionId !== currentFocus?.sessionId); + } + }, + newVisit() { + return createWorkHubRoutePolicyVisit(correctionMemory); + }, rememberTarget(target) { if (currentFocus?.sessionId === target.sessionId) return; previousFocus = currentFocus; currentFocus = target; }, - rememberCorrection(text, target) { - correctionSequence += 1; - corrections.unshift({ text, target, sequence: correctionSequence }); + reserveSubmissionOrder() { + correctionMemory.sequence += 1; + return correctionMemory.sequence; + }, + rememberCorrection(text, target, order) { + corrections.push({ text, target, sequence: order }); + corrections.sort((left, right) => right.sequence - left.sequence); corrections.splice(MAX_ROUTE_CORRECTIONS); }, }; } +function rankExactSessions( + text: string, + sessions: WorkHubSessionFacts[], +): Array<{ session: WorkHubSessionFacts; matchLength: number }> { + return sessions.map((session) => { + const qualifiedName = `${session.projectName}/${session.sessionName}`; + return { + session, + matchLength: Math.max( + exactIdentityMatchLength(text, qualifiedName), + exactIdentityMatchLength(text, session.sessionName), + ), + }; + }).filter(({ matchLength }) => matchLength >= MIN_EXACT_SESSION_NAME_LENGTH) + .sort((left, right) => right.matchLength - left.matchLength); +} + function correctedTarget( text: string, sessions: WorkHubSessionFacts[], @@ -297,6 +406,21 @@ function looksLikePreviousFocus(value: string): boolean { ); } +function naturalCorrectionTargetText(value: string): string | undefined { + const replacement = '(?:应该(?:是|用|改成|改为|切到|转到)?|而是|改成|改为|换成|换到|切到|转到|用|是)'; + const chinese = value.match( + new RegExp(`(?:(?:不是|不要再继续)\\s*(?:(?:这个|那个|当前这个|刚才那个)(?:工作|任务|Session|会话)?|[^,,。;;\\n]{1,32}(?:那个|那项工作|Session|会话|工作|任务))|(?:(?:这个|那个|当前这个|刚才那个)(?:工作|任务|Session|会话)|[^,,。;;\\n]{1,32}(?:那个|那项工作|Session|会话|工作|任务))\\s*(?:不对|搞错了|弄错了))[,,。;;\\n]\\s*${replacement}\\s*(.{2,})$`, 'iu'), + )?.[1]?.trim(); + if (chinese) return chinese; + return value.match( + /\b(?:not\s+(?:(?:this|that|the\s+current)(?:\s+(?:one|session|work|task))?|[^,.;\n]{1,32}\s+(?:session|work|task))|wrong\s+(?:one|session|work|task))\b[,.;\s]{0,4}(?:use|switch\s+to|change\s+to|move\s+to)\s+(.{2,})$/iu, + )?.[1]?.trim(); +} + +function looksLikeContentReplacement(value: string): boolean { + return /(?:不要(?:再)?用|别用|[^,,。;;\n]{1,32}(?:配置|实现|方案|字段|参数)?(?:不对|错了))[^\n]{0,64}[,,。;;]\s*(?:改成|改为|换成|换用|用)\s*\S{2,}/iu.test(value); +} + function looksLikeTargetUncertainty(value: string): boolean { return /(?:不确定(?:具体)?(?:是)?哪(?:一)?个|不知道(?:应该)?(?:选|继续|处理)哪(?:一)?个|可能是多个|哪个都可能|\b(?:i(?:'m| am)\s+)?not\s+sure\s+(?:which|where)|\b(?:i\s+)?(?:do\s+not|don't)\s+know\s+(?:which|where)|\b(?:could|might|may)\s+(?:be|belong\s+to)\s+(?:more\s+than\s+one|multiple)|\bwhich\s+(?:one|session|work|task)\b)/iu.test( value, diff --git a/apps/desktop/src/renderer/workhub-session-port.ts b/apps/desktop/src/renderer/workhub-session-port.ts index 47234d0b37..4d66344637 100644 --- a/apps/desktop/src/renderer/workhub-session-port.ts +++ b/apps/desktop/src/renderer/workhub-session-port.ts @@ -22,6 +22,7 @@ import type { DesktopTranscriptBatch, DesktopTranscriptHandle, } from '../preload/transcript-contract.js'; +import { parseDesktopSessionKey } from '../shared/runtime-host-identity.js'; import { DesktopTranscriptRangeStore } from './desktop-transcript-range-store.js'; import type { WorkHubProjectedTurn, @@ -30,7 +31,10 @@ import type { WorkHubSessionState, WorkHubSessionTarget, } from './workhub-controller.js'; -import { boundedWorkHubTimelineText } from './workhub-controller.js'; +import { + boundedWorkHubTimelineText, + WorkHubSessionSubmitError, +} from './workhub-controller.js'; export interface WorkHubDesktopSession { id: string; @@ -49,6 +53,10 @@ export interface WorkHubDesktopSession { export interface WorkHubDesktopSessionBridge { list(): Promise; + listWithCoverage?(): Promise<{ + sessions: readonly WorkHubDesktopSession[]; + completeHostIds: readonly string[]; + }>; listTurns(sessionId: string): Promise; create(input: { name: string }): Promise; send( @@ -104,15 +112,37 @@ export function createDesktopWorkHubSessionPort(deps: { : 'ordinary', archived: session.isArchived, state: projectState(session), + ...(session.runningTurnIds !== undefined + ? { runningTurnIds: [...session.runningTurnIds] } + : {}), ...(session.lastMessagePreview ? { latestResult: session.lastMessagePreview } : {}), updatedAt: session.lastMessageAt ?? session.statusUpdatedAt ?? 0, }); + const projectCatalog = async () => { + const snapshot = deps.sessions.listWithCoverage + ? await deps.sessions.listWithCoverage() + : { sessions: await deps.sessions.list(), completeHostIds: [] }; + const completeHostIds = new Set(snapshot.completeHostIds); + return { + sessions: snapshot.sessions.map(projectSession), + isCompleteFor(target: WorkHubSessionTarget) { + try { + return completeHostIds.has(parseDesktopSessionKey(target.sessionId).hostId); + } catch { + return false; + } + }, + }; + }; return { async list() { - return (await deps.sessions.list()).map(projectSession); + return (await projectCatalog()).sessions; + }, + listCatalog() { + return projectCatalog(); }, async recentTurns(targets) { const turnsBySession = await Promise.all( @@ -157,19 +187,60 @@ export function createDesktopWorkHubSessionPort(deps: { async create({ name }) { return projectSession(await deps.sessions.create({ name })); }, - async submit(target: WorkHubSessionTarget, text: string) { - const turnId = deps.newTurnId(); - const result = await deps.sessions.send(target.sessionId, { - type: 'send', - turnId, - text, - }); - if (!result.ok) throw new Error(`WorkHub Session send failed: ${result.reason}`); + reserveTurnId() { + return deps.newTurnId(); + }, + async submit(target: WorkHubSessionTarget, text: string, turnId: string) { + let result: Awaited>; + try { + result = await deps.sessions.send(target.sessionId, { + type: 'send', + turnId, + text, + }); + } catch (cause) { + throw new WorkHubSessionSubmitError( + 'WorkHub Session delivery outcome is unknown', + 'unknown', + { cause }, + ); + } + if (!result.ok) { + throw new WorkHubSessionSubmitError( + `WorkHub Session send failed: ${result.reason}`, + 'rejected', + ); + } return { turnId: result.turnId, ...(result.steered ? { steered: true as const } : {}), }; }, + async reconcileSubmission(target, reservedTurnId) { + try { + const messages = await readWorkHubSessionMessages(deps.transcripts, target); + let message: Extract | undefined; + for (let index = messages.length - 1; index >= 0; index -= 1) { + const entry = messages[index]; + if ( + entry?.type === 'user' && + (entry.turnId === reservedTurnId || entry.id === reservedTurnId) + ) { + message = entry; + break; + } + } + if (!message) return { kind: 'unknown' }; + if (message.turnId === reservedTurnId) { + return { kind: 'root', turnId: message.turnId }; + } + return message.steeringEventId + ? { kind: 'steered' } + : { kind: 'root', turnId: message.turnId }; + } catch { + return { kind: 'unknown' }; + } + }, async stop(target, expectedTurnId) { await deps.sessions.stop(target.sessionId, { source: 'stop_button', diff --git a/apps/desktop/src/renderer/workhub-surface.tsx b/apps/desktop/src/renderer/workhub-surface.tsx index 143241c780..2c755452dd 100644 --- a/apps/desktop/src/renderer/workhub-surface.tsx +++ b/apps/desktop/src/renderer/workhub-surface.tsx @@ -18,7 +18,12 @@ */ import { useCallback, useEffect, useRef, useState, type ReactNode } from 'react'; -import { ChatMessage, ChatMessageBubble, ChatMessageList } from '@astryxdesign/core'; +import { + ChatMessage, + ChatMessageBubble, + ChatMessageList, + Skeleton, +} from '@astryxdesign/core'; import { Button } from '@astryxdesign/core/Button'; import type { UiLocale } from '@maka/core/ui-locale'; import { ChatSurfaceLayout, Composer } from '@maka/ui'; @@ -147,38 +152,45 @@ export function projectedWorkHubTurnPresentation( export function WorkHubSurface(props: { controller: WorkHubController; locale: UiLocale; + initialFocusSessionId?: string; onOpenSession(sessionId: string): void; }) { const copy = workHubCopy(props.locale); const [projection, setProjection] = useState({ sessions: [], turns: [] }); const [turns, setTurns] = useState([]); const [pending, setPending] = useState(false); + const [initialLoadSettled, setInitialLoadSettled] = useState(false); // React state paints the lock; the gate closes the same-frame window before // a rerender can disable Composer and clarification controls. const routeGate = useRef(new WorkHubSurfaceRouteGate()).current; const refreshGate = useRef(new WorkHubProjectionRefreshGate()).current; const [loadError, setLoadError] = useState(false); - const refresh = useCallback(async () => { + const refresh = useCallback(async (focusSessionId?: string) => { const isLatest = refreshGate.begin(); try { - const next = await props.controller.read(); + const next = await props.controller.read(focusSessionId + ? { focus: { sessionId: focusSessionId } } + : undefined); if (!isLatest()) return; setProjection(next); setLoadError(false); + setInitialLoadSettled(true); } catch { if (!isLatest()) return; setLoadError(true); + setInitialLoadSettled(true); } }, [props.controller, refreshGate]); useEffect(() => { - void refresh(); + void refresh(props.initialFocusSessionId); const unsubscribe = props.controller.subscribe(() => void refresh()); return () => { refreshGate.invalidate(); unsubscribe(); + props.controller.resetVisitContext(); }; - }, [props.controller, refresh, refreshGate]); + }, [props.controller, props.initialFocusSessionId, refresh, refreshGate]); const route = useCallback(async ( input: WorkHubSubmitInput, @@ -217,14 +229,14 @@ export function WorkHubSurface(props: { const send = useCallback(async (value: string) => { const text = value.trim(); - if (!text || routeGate.pending) return false; + if (!text || !initialLoadSettled || routeGate.pending) return false; const requestId = crypto.randomUUID(); setTurns((current) => [...current, { requestId, text, state: 'routing' }]); const result = await route({ requestId, text }); // Composer clears only accepted drafts. Waiting, delivery failures, and a // ref-blocked duplicate keep the exact text available for retry. return workHubSubmissionClearsDraft(result); - }, [route, routeGate]); + }, [initialLoadSettled, route, routeGate]); const projectedTurns = visibleWorkHubProjectedTurns(projection.turns, turns); const conversationEmpty = projectedTurns.length === 0 && turns.length === 0; @@ -237,7 +249,7 @@ export function WorkHubSurface(props: { draftKey="workhub" onSend={send} onStop={() => {}} - sendBlocked={pending} + sendBlocked={pending || !initialLoadSettled} modelLabel="WorkHub" /> )} @@ -248,7 +260,9 @@ export function WorkHubSurface(props: {

WorkHub

{copy.subtitle}

- {copy.workCount(projection.sessions.length)} + {initialLoadSettled + ? copy.workCount(projection.sessions.length) + : copy.loading}
@@ -258,7 +272,9 @@ export function WorkHubSurface(props: { gap={4} isStreaming={pending} > - {loadError ? ( + {!initialLoadSettled ? ( + + ) : loadError ? (
{copy.loadFailed}
) : conversationEmpty ? (
@@ -287,6 +303,9 @@ export function WorkHubSurface(props: { requestId: turn.requestId, text: turn.text, explicitTarget: target, + ...(turn.outcome?.kind === 'clarification' && turn.outcome.correction + ? { correction: turn.outcome.correction } + : {}), })} onCorrect={(from, target) => void route({ requestId: turn.requestId, @@ -310,6 +329,27 @@ export function WorkHubSurface(props: { ); } +function WorkHubLoadingState(props: { label: string }) { + return ( +
+ {[0, 1].map((index) => ( + + ))} +
+ ); +} + function ProjectedWorkHubTurnView(props: { turn: WorkHubProjectedTurn; projection: WorkHubProjection; @@ -533,6 +573,7 @@ function workHubCopy(locale: UiLocale) { waitingForDecision: '这项工作正在等待你的决定。', requestNotSent: '新请求尚未发送;处理原 Session 中的交互后可以再次发送。', routing: '正在判断应该交给哪个 Session…', loadFailed: '无法读取已有工作。', + loading: '正在读取已有工作…', submitFailed: '输入未能送达,请重试。', scrollToBottom: '滚动到底部', archived: '已归档', states: { active: '活跃', running: '进行中', waiting_for_user: '等待你', blocked: '受阻', aborted: '已中止' }, turnStates: { running: '进行中', completed: '已完成', aborted: '已中止', failed: '失败' }, @@ -557,6 +598,7 @@ function workHubCopy(locale: UiLocale) { waitingForDecision: 'This work is waiting for your decision.', requestNotSent: 'The new request was not sent. Resolve the interaction in its Session, then send again.', routing: 'Choosing the right Session…', loadFailed: 'Could not read existing work.', + loading: 'Loading existing work…', submitFailed: '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' }, turnStates: { running: 'Running', completed: 'Completed', aborted: 'Aborted', failed: 'Failed' }, diff --git a/docs/astryx-surface-file-inventory.md b/docs/astryx-surface-file-inventory.md index 4707983699..834cd00c6d 100644 --- a/docs/astryx-surface-file-inventory.md +++ b/docs/astryx-surface-file-inventory.md @@ -5,7 +5,7 @@ Each row is one on-disk product surface file. Regenerated inventory must stay in Wiki bar: Design Conventions · API Use-the-System · Theming · Container Padding. -**Totals:** 209 files — blocker 0, polish 1, aligned 208. +**Totals:** 211 files — blocker 0, polish 1, aligned 210. ## Exclusions (explicit) diff --git a/docs/workhub-domain-language.md b/docs/workhub-domain-language.md index 3d0231fd80..ef1467b476 100644 --- a/docs/workhub-domain-language.md +++ b/docs/workhub-domain-language.md @@ -19,18 +19,57 @@ # WorkHub domain language -WorkHub gives users one conversational place to continue, create, and inspect work while ordinary Sessions remain the product's canonical work records. +WorkHub gives users one persistent conversational place to ask, clarify, continue, +create, and inspect work. It is backed by one stable Coordination Session per +Runtime Host while concrete execution remains authoritative in ordinary Sessions. + +This document names the approved target architecture. The current R2.4 +implementation is a transitional deterministic router and does not yet create the +Coordination Session described below. ## Terms -**Session**: The authoritative record for identity, transcript, execution state, permissions, interactions, and recovery. A Session ID is the stable identity of the work. +**Session**: The existing transcript, execution-boundary, permission, interaction, +and recovery substrate. A Session owns only the conversation or execution admitted +to that Session. + +**Coordination Session**: A special Session role used by WorkHub. Each Runtime Host +has at most one stable Coordination Session. It owns WorkHub user messages, ordinary +Q&A, clarification, coordination decisions, delegation references, and coordination +summaries. It is hidden from the ordinary Session list and never routes to itself. -**Work**: The user-facing continuity of exactly one ordinary Session. “Work” is a product-language view of a Session, not a second stored record. +**Ordinary Session**: A Session that owns concrete work execution, including its +project/filesystem scope, model and permissions, root-Turn admission, tools, +artifacts, recovery, lifecycle, and authoritative execution transcript. -**WorkHub**: A projection and routing surface over ordinary Sessions. It may keep transient inference context while mounted, but it does not own a transcript or execution state. +**Work**: User-facing continuity around a goal. Whether Work is 1:1 with Session, +1:N over Sessions, or an independent durable entity is deliberately unresolved. + +**WorkHub**: The unified conversational entry and coordination surface backed by +the active Runtime Host's Coordination Session. It may answer locally, clarify, +delegate to an existing ordinary Session, or create a new ordinary Session. **Session projection**: A rebuildable view derived from Session facts for display and routing. It can be discarded and recreated without losing work. -**Route correction**: A user's decision that an input belongs to a different existing Session. It may influence later transient routing, but it does not become an authority for Session content or state. +**Disposition**: The coordination outcome for one WorkHub input: +`answer_here`, `delegate_existing`, `create_new`, or `clarify`. + +**Delegation**: A reference from a Coordination Turn to a target ordinary Session +and Turn. Delegation links the two authoritative transcripts; it does not copy the +target execution transcript into WorkHub. + +**Action Gate**: The deterministic Runtime boundary that validates a proposed +disposition, target, creation, Stop, confirmation, tools, and permissions. A model +or routing policy may propose an action but cannot authorize it. + +**Route correction**: A user's decision that an input belongs to a different +existing Session. R2.4 retains only bounded inference memory for later target +resolution. Correction precedence follows user submission order, not asynchronous +completion order, and it never replaces either Session's transcript authority. + +**R2.4**: The deterministic context-continuity routing baseline. It remains useful +as an experiment baseline or target resolver behind WorkHub's coordination layer; +it is not the final architecture or authority boundary of WorkHub. -_Avoid_: independent Work records, copied transcripts, or a second writable WorkHub state store. +_Avoid_: copied execution transcripts, self-routing, a second Session/WorkHub +storage substrate, or treating model/routing output as execution authority.