diff --git a/mobile/src/components/PickerModal.tsx b/mobile/src/components/PickerModal.tsx index 9456307091f..45b6593c09f 100644 --- a/mobile/src/components/PickerModal.tsx +++ b/mobile/src/components/PickerModal.tsx @@ -15,6 +15,7 @@ export type PickerOption = { type Props = { visible: boolean title: string + subtitle?: string options: PickerOption[] selected: T onSelect: (value: T) => void @@ -32,6 +33,7 @@ type PickerModalContentProps = Pick< export function PickerModal({ visible, title, + subtitle, options, selected, onSelect, @@ -44,6 +46,7 @@ export function PickerModal({ {title} + {subtitle ? {subtitle} : null} state.setShowFilterModal(false)}> - Filter + + Filter + {WORKSPACE_VIEW_SHARED_NOTE} + {settings.activeFilterCount > 0 && ( Clear filters diff --git a/mobile/src/host-screen/host-screen-secondary-styles.ts b/mobile/src/host-screen/host-screen-secondary-styles.ts index af14aa2b97a..7f6c4f3f65c 100644 --- a/mobile/src/host-screen/host-screen-secondary-styles.ts +++ b/mobile/src/host-screen/host-screen-secondary-styles.ts @@ -47,11 +47,19 @@ export const hostScreenSecondaryStyles = StyleSheet.create({ paddingHorizontal: spacing.xs, marginBottom: spacing.md }, + filterModalHeading: { + flexShrink: 1 + }, filterModalTitle: { fontSize: 15, fontWeight: '600', color: colors.textPrimary }, + filterModalSubtitle: { + fontSize: 11, + color: colors.textMuted, + marginTop: 2 + }, clearFiltersText: { fontSize: 13, color: colors.textSecondary diff --git a/mobile/src/session/use-mobile-structured-agent-session.ts b/mobile/src/session/use-mobile-structured-agent-session.ts index 653d7fbaaa1..13c33050b93 100644 --- a/mobile/src/session/use-mobile-structured-agent-session.ts +++ b/mobile/src/session/use-mobile-structured-agent-session.ts @@ -158,8 +158,13 @@ export function useMobileStructuredAgentSession(args: { }) const messages = useMemo( - () => projectStructuredAgentSessionMessages(state.items, [], state.submissions), - [state.items, state.submissions] + // A message the host recorded and then rejected stays in place, said once on its row. + () => + projectStructuredAgentSessionMessages(state.items, [], state.submissions, { + rejectedInPlace: true, + queuedMessageIds: (queuedMessages ?? []).map((draft) => draft.messageId) + }), + [queuedMessages, state.items, state.submissions] ) const turnId = activeStructuredAgentSessionTurnId(state.items) const turnTiming = useMobileStructuredAgentTurnTiming(state, turnId) diff --git a/mobile/src/worktree/workspace-list-picker-options.ts b/mobile/src/worktree/workspace-list-picker-options.ts index 180d4038320..ca597ed38cc 100644 --- a/mobile/src/worktree/workspace-list-picker-options.ts +++ b/mobile/src/worktree/workspace-list-picker-options.ts @@ -1,6 +1,9 @@ import type { PickerOption } from '../components/PickerModal' import type { MobileGroupMode, MobileSortMode } from './workspace-view-settings' +// Why: the host may be headless, so the note can't promise a desktop sidebar. +export const WORKSPACE_VIEW_SHARED_NOTE = 'Synced across your devices' + export const WORKSPACE_SORT_OPTIONS: PickerOption[] = [ // Why: desktop and persisted state keep the `smart` key, while mobile shows the product label. { @@ -11,7 +14,7 @@ export const WORKSPACE_SORT_OPTIONS: PickerOption[] = [ { value: 'name', label: 'Name', subtitle: 'Alphabetical by name' }, { value: 'recent', label: 'Recent', subtitle: 'Most recent output first' }, { value: 'repo', label: 'Repo', subtitle: 'Repository, then workspace name' }, - { value: 'manual', label: 'Manual', subtitle: 'Server order' } + { value: 'manual', label: 'Manual', subtitle: 'Desktop drag order' } ] export const WORKSPACE_GROUP_OPTIONS: PickerOption[] = [ diff --git a/mobile/src/worktree/workspace-view-settings.ts b/mobile/src/worktree/workspace-view-settings.ts index f1d17189adc..221fa7211be 100644 --- a/mobile/src/worktree/workspace-view-settings.ts +++ b/mobile/src/worktree/workspace-view-settings.ts @@ -7,7 +7,7 @@ import type { WorkspaceStatusDefinition } from '../../../src/shared/worktree/typ import { coerceMobileWorkspaceStatuses } from './mobile-workspace-statuses' export type MobileGroupMode = 'none' | 'workspaceStatus' | 'repo' | 'prStatus' -// Desktop sort adds 'manual'; mobile renders it but sorts by server order. +// Desktop sort adds 'manual'; mobile orders it by the desktop's drag ranks. export type MobileSortMode = 'smart' | 'name' | 'recent' | 'repo' | 'manual' // Desktop PersistedUIState fields this screen syncs (a structural subset). diff --git a/src/main/ai-vault/session-scanner-opencode-cancellation.test.ts b/src/main/ai-vault/session-scanner-opencode-cancellation.test.ts index 4f47af380cd..084a73cd661 100644 --- a/src/main/ai-vault/session-scanner-opencode-cancellation.test.ts +++ b/src/main/ai-vault/session-scanner-opencode-cancellation.test.ts @@ -1,4 +1,7 @@ import { afterEach, beforeEach, describe, expect, it, vi } from 'vitest' +import { mkdtempSync, rmSync, writeFileSync } from 'node:fs' +import { tmpdir } from 'node:os' +import { join } from 'node:path' import type * as workerSpawn from './session-scanner-opencode-sqlite-worker-spawn' import type { SessionFileDiscovery } from './session-scanner-types' import type { TranscriptReadOutcome } from './session-transcript-consumers' @@ -24,6 +27,7 @@ import { registerTranscriptConsumer, resetTranscriptConsumersForTests } from './session-transcript-consumers' +import { runOpenCodeSqliteScanRequest } from './session-scanner-opencode-sqlite-scan-scope' const file = { path: '/fixture/opencode.db#session', @@ -40,6 +44,7 @@ beforeEach(() => { resetSessionParseCacheForTests() }) afterEach(() => { + vi.useRealTimers() resetTranscriptConsumersForTests() resetSessionParseCacheForTests() }) @@ -48,7 +53,11 @@ function configure(agent: 'opencode' | 'opencode2') { readers.discover.mockResolvedValue([{ agent, rootDir: '/fixture', files: [file] }]) const accumulator = createAccumulator({ agent, file, sessionId: 'session' }) accumulator.title = 'SQLite session' - return finalizeSession(accumulator, 'linux') + const session = finalizeSession(accumulator, 'linux') + if (!session) { + throw new Error('Configured SQLite session was empty') + } + return session } function untilAborted(signal: AbortSignal | undefined): Promise { @@ -61,6 +70,83 @@ function untilAborted(signal: AbortSignal | undefined): Promise { } describe.each(['opencode', 'opencode2'] as const)('%s scan cancellation', (agent) => { + it('reports a deadline while retaining completed sessions and other agents, then retries', async () => { + vi.useFakeTimers() + const root = mkdtempSync(join(tmpdir(), 'orca-scan-deadline-')) + try { + const session = configure(agent) + const blocked = { ...file, path: '/fixture/opencode.db#blocked' } + const claudePath = join(root, 'claude.jsonl') + writeFileSync( + claudePath, + `${JSON.stringify({ + type: 'user', + sessionId: 'retained-claude', + timestamp: '2026-05-01T10:00:00.000Z', + cwd: root, + message: { role: 'user', content: 'Retain this other-agent session' } + })}\n` + ) + readers.discover.mockResolvedValue([ + { agent, rootDir: '/fixture', files: [file, blocked] }, + { agent: 'claude', rootDir: root, files: [{ ...file, path: claudePath }] } + ]) + readers.parse.mockImplementation(({ sessionId, signal }) => + sessionId === 'blocked' + ? runOpenCodeSqliteScanRequest(signal, untilAborted) + : Promise.resolve(session) + ) + const pending = scanAiVaultSessions({ platform: 'linux' }) + await vi.waitFor(() => expect(readers.parse).toHaveBeenCalledTimes(2)) + await vi.advanceTimersByTimeAsync(45_000) + const result = await pending + expect(result.sessions.map((row) => row.sessionId)).toEqual( + expect.arrayContaining(['session', 'retained-claude']) + ) + expect(result.sessions).toHaveLength(2) + expect(result.issues).toEqual([ + expect.objectContaining({ + agent, + path: blocked.path, + message: expect.stringContaining('45s work budget') + }) + ]) + readers.parse.mockResolvedValue({ ...session, sessionId: 'blocked', filePath: blocked.path }) + const recovered = await scanAiVaultSessions({ platform: 'linux' }) + expect(recovered.sessions).toHaveLength(3) + expect( + readers.parse.mock.calls.filter(([args]) => args.sessionId === 'blocked') + ).toHaveLength(2) + } finally { + rmSync(root, { recursive: true, force: true }) + } + }) + + it('marks a deadline capture incomplete instead of caching a failed history', async () => { + vi.useFakeTimers() + const session = configure(agent) + const outcomes: TranscriptReadOutcome[] = [] + registerTranscriptConsumer({ + beginRead: () => ({ message() {}, finish: (outcome) => outcomes.push(outcome) }) + }) + readers.capture.mockImplementationOnce(({ signal }) => + runOpenCodeSqliteScanRequest(signal, untilAborted) + ) + const pending = scanAiVaultSessions({ platform: 'linux' }) + await vi.waitFor(() => expect(readers.capture).toHaveBeenCalledOnce()) + await vi.advanceTimersByTimeAsync(45_000) + const result = await pending + expect(result.sessions).toEqual([]) + expect(result.issues).toEqual([ + expect.objectContaining({ message: expect.stringContaining('45s work budget') }) + ]) + expect(outcomes).toEqual([{ session: null, byteOffset: 0, incomplete: true }]) + readers.capture.mockResolvedValue({ session, messages }) + expect((await scanAiVaultSessions({ platform: 'linux' })).sessions).toHaveLength(1) + expect(readers.capture).toHaveBeenCalledTimes(2) + expect(outcomes.at(-1)?.incomplete).toBe(false) + }) + it.each(['parse', 'capture'] as const)( 'cancels an active %s and retries the uncached read', async (mode) => { diff --git a/src/main/ai-vault/session-scanner-opencode-sqlite-scan-scope.test.ts b/src/main/ai-vault/session-scanner-opencode-sqlite-scan-scope.test.ts new file mode 100644 index 00000000000..bc97f1aa388 --- /dev/null +++ b/src/main/ai-vault/session-scanner-opencode-sqlite-scan-scope.test.ts @@ -0,0 +1,138 @@ +import { afterEach, expect, it, vi } from 'vitest' +import { + OPENCODE_SQLITE_SCAN_BUDGET_MS, + runOpenCodeSqliteScanRequest, + withOpenCodeSqliteScanScope +} from './session-scanner-opencode-sqlite-scan-scope' + +afterEach(() => vi.useRealTimers()) + +function waitForAbort(signal: AbortSignal | undefined): Promise { + if (!signal) { + throw new Error('Missing scoped request signal') + } + return new Promise((_resolve, reject) => { + signal.addEventListener('abort', () => reject(signal.reason), { once: true }) + }) +} + +function wait(ms: number): Promise { + return new Promise((resolve) => setTimeout(resolve, ms)) +} + +it('spends the budget while admission is pending and refuses later work in that scan', async () => { + vi.useFakeTimers() + const admitted = vi.fn() + const outcome = withOpenCodeSqliteScanScope(async () => { + const error = await runOpenCodeSqliteScanRequest(undefined, waitForAbort).catch((err) => err) + expect(error).toMatchObject({ name: 'OpenCodeSqliteScanDeadlineError' }) + await expect(runOpenCodeSqliteScanRequest(undefined, admitted)).rejects.toBe(error) + }) + await vi.advanceTimersByTimeAsync(OPENCODE_SQLITE_SCAN_BUDGET_MS) + await outcome + expect(admitted).not.toHaveBeenCalled() + expect(vi.getTimerCount()).toBe(0) +}) + +it('banks only outstanding work across legs and does not spend other-agent time', async () => { + vi.useFakeTimers() + const outcome = withOpenCodeSqliteScanScope(async () => { + await runOpenCodeSqliteScanRequest(undefined, () => wait(20_000)) + await wait(70_000) + return runOpenCodeSqliteScanRequest(undefined, waitForAbort) + }).catch((error) => error) + await vi.advanceTimersByTimeAsync(90_000) + let completed = false + void outcome.then(() => { + completed = true + }) + await vi.advanceTimersByTimeAsync(24_999) + expect(completed).toBe(false) + await vi.advanceTimersByTimeAsync(1) + expect(await outcome).toMatchObject({ name: 'OpenCodeSqliteScanDeadlineError' }) + expect(vi.getTimerCount()).toBe(0) +}) + +it('counts overlapping preparation and worker waits once', async () => { + vi.useFakeTimers() + const outcome = withOpenCodeSqliteScanScope(() => + runOpenCodeSqliteScanRequest(undefined, () => + Promise.all([ + runOpenCodeSqliteScanRequest(undefined, waitForAbort), + runOpenCodeSqliteScanRequest(undefined, waitForAbort) + ]) + ) + ).catch((error) => error) + await vi.advanceTimersByTimeAsync(OPENCODE_SQLITE_SCAN_BUDGET_MS - 1) + expect(vi.getTimerCount()).toBe(1) + await vi.advanceTimersByTimeAsync(1) + expect(await outcome).toMatchObject({ name: 'OpenCodeSqliteScanDeadlineError' }) + expect(vi.getTimerCount()).toBe(0) +}) + +it('keeps concurrent scans independent and gives the next scan a fresh owner and budget', async () => { + vi.useFakeTimers() + const owners: unknown[] = [] + const first = withOpenCodeSqliteScanScope(() => + runOpenCodeSqliteScanRequest(undefined, (signal, owner) => { + owners.push(owner) + return waitForAbort(signal) + }) + ).catch((error) => error) + await vi.advanceTimersByTimeAsync(30_000) + const second = withOpenCodeSqliteScanScope(() => + runOpenCodeSqliteScanRequest(undefined, (signal, owner) => { + owners.push(owner) + return waitForAbort(signal) + }) + ).catch((error) => error) + await vi.advanceTimersByTimeAsync(15_000) + expect(await first).toMatchObject({ name: 'OpenCodeSqliteScanDeadlineError' }) + expect(vi.getTimerCount()).toBe(1) + await vi.advanceTimersByTimeAsync(30_000) + expect(await second).toMatchObject({ name: 'OpenCodeSqliteScanDeadlineError' }) + await withOpenCodeSqliteScanScope(() => + runOpenCodeSqliteScanRequest(undefined, async (_signal, owner) => { + owners.push(owner) + }) + ) + expect(new Set(owners).size).toBe(3) + expect(vi.getTimerCount()).toBe(0) +}) + +it('preserves the caller cancellation reason and disposes its timer', async () => { + vi.useFakeTimers() + const controller = new AbortController() + const reason = new Error('caller cancelled') + const outcome = withOpenCodeSqliteScanScope(() => + runOpenCodeSqliteScanRequest(controller.signal, waitForAbort) + ).catch((error) => error) + controller.abort(reason) + expect(await outcome).toBe(reason) + expect(vi.getTimerCount()).toBe(0) +}) + +it('leaves unrelated native-chat and Zcode calls unscoped and retires the scan signal', async () => { + vi.useFakeTimers() + let scopedSignal: AbortSignal | undefined + const caller = new AbortController() + await withOpenCodeSqliteScanScope(async () => { + for (const agent of ['zcode', 'native-chat'] as const) { + await runOpenCodeSqliteScanRequest( + caller.signal, + async (signal, owner) => { + expect(signal).toBe(caller.signal) + expect(owner).toBeUndefined() + expect(vi.getTimerCount()).toBe(0) + }, + agent + ) + } + await runOpenCodeSqliteScanRequest(undefined, async (signal) => { + scopedSignal = signal + }) + }) + expect(scopedSignal?.aborted).toBe(true) + expect(caller.signal.aborted).toBe(false) + expect(vi.getTimerCount()).toBe(0) +}) diff --git a/src/main/ai-vault/session-scanner-opencode-sqlite-scan-scope.ts b/src/main/ai-vault/session-scanner-opencode-sqlite-scan-scope.ts new file mode 100644 index 00000000000..e3a53a535b2 --- /dev/null +++ b/src/main/ai-vault/session-scanner-opencode-sqlite-scan-scope.ts @@ -0,0 +1,75 @@ +import { AsyncLocalStorage } from 'node:async_hooks' +import { throwIfSignalAborted } from '../../shared/abort-signal-reason' +import type { WorkerThreadRequestOwner } from '../worker-thread-request-queue' + +export const OPENCODE_SQLITE_SCAN_BUDGET_MS = 45_000 + +class OpenCodeSqliteScanScope implements WorkerThreadRequestOwner { + private readonly controller = new AbortController() + readonly signal = this.controller.signal + private remainingMs = OPENCODE_SQLITE_SCAN_BUDGET_MS + private outstanding = 0 + private armedAt = 0 + private timer: NodeJS.Timeout | undefined + + async run( + callerSignal: AbortSignal | undefined, + fn: (signal: AbortSignal, owner: WorkerThreadRequestOwner) => Promise + ): Promise { + throwIfSignalAborted(callerSignal) + throwIfSignalAborted(this.signal) + if (this.outstanding++ === 0) { + this.armedAt = Date.now() + this.timer = setTimeout(() => { + const error = new Error( + `OpenCode SQLite scan exceeded its ${OPENCODE_SQLITE_SCAN_BUDGET_MS / 1000}s work budget` + ) + error.name = 'OpenCodeSqliteScanDeadlineError' + this.controller.abort(error) + }, this.remainingMs) + this.timer.unref?.() + } + const signal = callerSignal ? AbortSignal.any([callerSignal, this.signal]) : this.signal + try { + return await fn(signal, this) + } finally { + if (--this.outstanding === 0) { + this.pause() + } + } + } + + dispose(): void { + this.pause() + this.controller.abort(new Error('OpenCode SQLite scan ended')) + } + + private pause(): void { + if (this.timer) { + clearTimeout(this.timer) + this.timer = undefined + this.remainingMs = Math.max(0, this.remainingMs - (Date.now() - this.armedAt)) + } + } +} + +const scanScope = new AsyncLocalStorage() + +export async function withOpenCodeSqliteScanScope(fn: () => Promise): Promise { + const scope = new OpenCodeSqliteScanScope() + try { + return await scanScope.run(scope, fn) + } finally { + scope.dispose() + } +} + +// The clock covers outstanding SQLite work, including admission and WSL preparation. +export function runOpenCodeSqliteScanRequest( + signal: AbortSignal | undefined, + fn: (signal?: AbortSignal, owner?: WorkerThreadRequestOwner) => Promise, + agent?: 'opencode2' | 'zcode' | 'native-chat' +): Promise { + const scope = agent === 'zcode' || agent === 'native-chat' ? undefined : scanScope.getStore() + return scope ? scope.run(signal, fn) : fn(signal) +} diff --git a/src/main/ai-vault/session-scanner-opencode-sqlite-worker-client.test.ts b/src/main/ai-vault/session-scanner-opencode-sqlite-worker-client.test.ts index 93df4ed987c..68fb8d09d90 100644 --- a/src/main/ai-vault/session-scanner-opencode-sqlite-worker-client.test.ts +++ b/src/main/ai-vault/session-scanner-opencode-sqlite-worker-client.test.ts @@ -12,6 +12,7 @@ import type { OpenCodeSqliteWorkerResponse } from './session-scanner-opencode-sqlite-worker-protocol' import type { AiVaultScanIssue } from '../../shared/ai-vault-types' +import { withOpenCodeSqliteScanScope } from './session-scanner-opencode-sqlite-scan-scope' // A worker_threads stand-in the tests drive directly: it records posted requests // and lets a test emit message/error/exit without a built worker bundle. @@ -75,6 +76,63 @@ function makeFactory(workers: FakeWorker[]): () => Worker { } describe('OpenCodeSqliteWorkerClient', () => { + it('expires a scan in the FIFO without cancelling ordinary reads or a later scan', async () => { + vi.useFakeTimers() + const workers: FakeWorker[] = [] + const client = new OpenCodeSqliteWorkerClient({ workerFactory: makeFactory(workers), log() {} }) + const issues: AiVaultScanIssue[] = [] + try { + const ordinary = [1, 2].map((id) => + client.list({ + dbPaths: [`/ordinary-${id}.db`], + limit: 1, + issues: [] + }) + ) + const scan = withOpenCodeSqliteScanScope(() => + client.list({ + dbPaths: ['/scan.db'], + limit: 1, + issues + }) + ) + await vi.advanceTimersByTimeAsync(25_000) + workers[0].emit('message', { + id: workers[0].lastId(), + ok: true, + value: { candidates: [], issues: [] } + }) + await vi.advanceTimersByTimeAsync(20_000) + await expect(scan).resolves.toEqual([]) + expect(issues).toEqual([ + expect.objectContaining({ + kind: 'scope', + message: expect.stringContaining('45s work budget') + }) + ]) + expect(workers[0].postedRequests).toHaveLength(2) + expect(workers[0].terminated).toBe(false) + workers[0].emit('message', { + id: workers[0].lastId(), + ok: true, + value: { candidates: [], issues: [] } + }) + await expect(Promise.all(ordinary)).resolves.toEqual([[], []]) + const next = withOpenCodeSqliteScanScope(() => + client.list({ dbPaths: ['/next.db'], limit: 1, issues: [] }) + ) + workers[0].emit('message', { + id: workers[0].lastId(), + ok: true, + value: { candidates: [], issues: [] } + }) + await expect(next).resolves.toEqual([]) + } finally { + client.dispose() + vi.useRealTimers() + } + }) + it('correlates responses by id and ignores stale ids', async () => { const workers: FakeWorker[] = [] const client = new OpenCodeSqliteWorkerClient({ workerFactory: makeFactory(workers), log() {} }) @@ -188,7 +246,7 @@ describe('OpenCodeSqliteWorkerClient', () => { client.list({ dbPaths: ['/tmp/opencode.db'], limit: 10, issues: listIssues }) ).resolves.toEqual([]) expect( - listIssues.some((issue) => /background scanner could not start/.test(issue.message)) + listIssues.some((issue) => issue.message.includes('background scanner could not start')) ).toBe(true) await expect( client.parse({ dbPath: '/tmp/opencode.db', sessionId: 'ses_skipped', platform: 'darwin' }) @@ -264,7 +322,7 @@ describe('OpenCodeSqliteWorkerClient', () => { const first = await client.list({ dbPaths: ['/db'], limit: 10, issues: firstIssues }) expect(first).toEqual([]) expect( - firstIssues.some((issue) => /background scanner could not start/.test(issue.message)) + firstIssues.some((issue) => issue.message.includes('background scanner could not start')) ).toBe(true) await expect( client.parse({ dbPath: '/db', sessionId: 'ses_heal', platform: 'darwin' }) diff --git a/src/main/ai-vault/session-scanner-opencode-sqlite-worker-client.ts b/src/main/ai-vault/session-scanner-opencode-sqlite-worker-client.ts index f303de52443..c6a85cb58b0 100644 --- a/src/main/ai-vault/session-scanner-opencode-sqlite-worker-client.ts +++ b/src/main/ai-vault/session-scanner-opencode-sqlite-worker-client.ts @@ -12,6 +12,7 @@ import type { import { parseOpenCodeSqliteCaptureValue } from './session-scanner-opencode-sqlite-worker-response' import type { SessionFileCandidate } from './session-scanner-types' import { errorMessage } from './session-scanner-values' +import { runOpenCodeSqliteScanRequest } from './session-scanner-opencode-sqlite-scan-scope' // Why (#8864): a lazily-spawned, unref'd worker runs OpenCode SQLite reads off // the main-process event loop. This module owns only the OpenCode legs; the @@ -117,7 +118,8 @@ export class OpenCodeSqliteWorkerClient { ...(args.agent ? { agent: args.agent } : {}) }), LIST_TIMEOUT_MS, - args.signal + args.signal, + args.agent )) as OpenCodeSqliteListValue args.issues.push(...value.issues) return value.candidates @@ -178,7 +180,8 @@ export class OpenCodeSqliteWorkerClient { ...(args.agent ? { agent: args.agent } : {}) }), PARSE_TIMEOUT_MS, - args.signal + args.signal, + args.agent ) // oxlint-disable-next-line typescript/consistent-type-assertions -- SAFETY: the worker's parse leg returns exactly this, built by the repo's own reader on the other side of a structured clone. return value as AiVaultSession | null @@ -217,7 +220,8 @@ export class OpenCodeSqliteWorkerClient { ...(args.agent ? { agent: args.agent } : {}) }), CAPTURE_TIMEOUT_MS, - args.signal + args.signal, + args.agent ) return parseOpenCodeSqliteCaptureValue(value) } catch (err) { @@ -229,7 +233,12 @@ export class OpenCodeSqliteWorkerClient { args: Omit, signal?: AbortSignal ): Promise { - const value = await this.dispatch((id) => ({ ...args, id }), PARSE_TIMEOUT_MS, signal) + const value = await this.dispatch( + (id) => ({ ...args, id }), + PARSE_TIMEOUT_MS, + signal, + 'native-chat' + ) // oxlint-disable-next-line typescript/consistent-type-assertions -- SAFETY: Only this build's internal worker dispatch constructs page/signal results; they are not client-supplied paths or frames. return value as OpenCodeNativeChatReadValue } @@ -241,13 +250,20 @@ export class OpenCodeSqliteWorkerClient { private async dispatch( buildRequest: (id: number) => OpenCodeSqliteWorkerRequest, timeoutMs: number, - signal?: AbortSignal + signal?: AbortSignal, + agent?: 'opencode2' | 'zcode' | 'native-chat' ): Promise { const deadline = this.requestTimeoutMs ?? timeoutMs - const response = await this.requests.dispatch( - (id) => ({ ...buildRequest(id), timeoutMs: deadline }), - deadline, - signal + const response = await runOpenCodeSqliteScanRequest( + signal, + (requestSignal, owner) => + this.requests.dispatch( + (id) => ({ ...buildRequest(id), timeoutMs: deadline }), + deadline, + requestSignal, + owner + ), + agent ) if (!response.ok) { throw new Error(response.error) diff --git a/src/main/ai-vault/session-scanner-opencode-sqlite-worker-spawn.ts b/src/main/ai-vault/session-scanner-opencode-sqlite-worker-spawn.ts index 13b9cd3cd5e..e89a6361b1e 100644 --- a/src/main/ai-vault/session-scanner-opencode-sqlite-worker-spawn.ts +++ b/src/main/ai-vault/session-scanner-opencode-sqlite-worker-spawn.ts @@ -16,6 +16,7 @@ import { openCodeWslPath } from './session-scanner-opencode-wsl-client' import { findForeignSqliteReaderEntry } from '../foreign-sqlite-readers/foreign-sqlite-reader-entry-path' +import { runOpenCodeSqliteScanRequest } from './session-scanner-opencode-sqlite-scan-scope' // Why: resolve the built worker entry + own the process-wide shared client so // the client class stays free of Electron (require'd lazily here) and the @@ -179,8 +180,14 @@ async function listForHost( const first = paths.values().next().value! const issues: AiVaultScanIssue[] = [] try { - const client = await openCodeWslClient(distro, first, args.signal) - const result = await client.list({ ...args, dbPaths: [...paths.keys()], issues }) + const result = await runOpenCodeSqliteScanRequest( + args.signal, + async (signal) => { + const client = await openCodeWslClient(distro, first, signal) + return client.list({ ...args, signal, dbPaths: [...paths.keys()], issues }) + }, + args.agent + ) return result.flatMap((candidate) => { const parsed = splitOpenCodeSqliteCandidate(candidate.file.path, args.agent) const original = parsed && paths.get(parsed.dbPath) @@ -223,8 +230,14 @@ async function parseForHost( if (!wsl) { return getSharedClient().parse(args) } - const client = await openCodeWslClient(wsl.distro, args.dbPath, args.signal) - const session = await client.parse({ ...args, dbPath: wsl.linuxPath, platform: 'linux' }) + const session = await runOpenCodeSqliteScanRequest( + args.signal, + async (signal) => { + const client = await openCodeWslClient(wsl.distro, args.dbPath, signal) + return client.parse({ ...args, signal, dbPath: wsl.linuxPath, platform: 'linux' }) + }, + args.agent + ) return mapOpenCodeWslSession(session, args.dbPath) } @@ -235,8 +248,14 @@ async function captureForHost( if (!wsl) { return getSharedClient().capture(args) } - const client = await openCodeWslClient(wsl.distro, args.dbPath, args.signal) - const capture = await client.capture({ ...args, dbPath: wsl.linuxPath, platform: 'linux' }) + const capture = await runOpenCodeSqliteScanRequest( + args.signal, + async (signal) => { + const client = await openCodeWslClient(wsl.distro, args.dbPath, signal) + return client.capture({ ...args, signal, dbPath: wsl.linuxPath, platform: 'linux' }) + }, + args.agent + ) return { ...capture, session: mapOpenCodeWslSession(capture.session, args.dbPath) } } diff --git a/src/main/ai-vault/session-scanner-opencode-wsl-routing.test.ts b/src/main/ai-vault/session-scanner-opencode-wsl-routing.test.ts index 754211c323c..9db4ba85e83 100644 --- a/src/main/ai-vault/session-scanner-opencode-wsl-routing.test.ts +++ b/src/main/ai-vault/session-scanner-opencode-wsl-routing.test.ts @@ -3,6 +3,7 @@ import type { AiVaultScanIssue } from '../../shared/ai-vault-types' import { createAccumulator, finalizeSession } from './session-scanner-accumulator' import type { SessionFileCandidate } from './session-scanner-types' import type * as wslClientModule from './session-scanner-opencode-wsl-client' +import { withOpenCodeSqliteScanScope } from './session-scanner-opencode-sqlite-scan-scope' const mocks = vi.hoisted(() => ({ native: { @@ -59,9 +60,53 @@ beforeEach(() => { vi.clearAllMocks() vi.spyOn(process, 'platform', 'get').mockReturnValue('win32') }) -afterEach(() => vi.restoreAllMocks()) +afterEach(() => { + vi.useRealTimers() + vi.restoreAllMocks() +}) describe('OpenCode SQLite execution-host routes', () => { + it('budgets WSL preparation while retaining a native source that already answered', async () => { + vi.useFakeTimers() + mocks.native.list.mockResolvedValueOnce([row(native)]) + mocks.guest.mockImplementationOnce((_distro, _path, signal: AbortSignal | undefined) => { + if (!signal) { + throw new Error('Missing scoped preparation signal') + } + return new Promise((_resolve, reject) => + signal.addEventListener('abort', () => reject(signal.reason), { once: true }) + ) + }) + const issues: AiVaultScanIssue[] = [] + const result = withOpenCodeSqliteScanScope(() => + listOpenCodeSqliteSessionsViaWorker({ + dbPaths: [native, ubuntu], + limit: 2, + issues + }) + ) + await vi.advanceTimersByTimeAsync(45_000) + expect((await result).map((entry) => entry.file.path)).toEqual([`${native}#same-session`]) + expect(issues).toEqual([ + expect.objectContaining({ + path: ubuntu, + kind: 'scope', + message: expect.stringContaining('45s work budget') + }) + ]) + mocks.guest.mockResolvedValueOnce({ list: vi.fn(async () => [row(guest)]) }) + const recoveredIssues: AiVaultScanIssue[] = [] + const recovered = await withOpenCodeSqliteScanScope(() => + listOpenCodeSqliteSessionsViaWorker({ + dbPaths: [ubuntu], + limit: 2, + issues: recoveredIssues + }) + ) + expect(recovered).toHaveLength(1) + expect(recoveredIssues).toEqual([]) + }) + it('separates native and distro databases and preserves equal IDs in different distros', async () => { const list = vi.fn(async (args) => [row(args.dbPaths[0])]) mocks.guest.mockResolvedValue({ list }) diff --git a/src/main/ai-vault/session-scanner.ts b/src/main/ai-vault/session-scanner.ts index 8716440a565..f0525c918c8 100644 --- a/src/main/ai-vault/session-scanner.ts +++ b/src/main/ai-vault/session-scanner.ts @@ -45,6 +45,7 @@ import { clampPositiveInteger, errorMessage } from './session-scanner-values' import { throwIfAiVaultScanCancelled } from './ai-vault-scan-cancellation' import { DEFAULT_AI_VAULT_SCAN_LIMIT } from '../../shared/ai-vault-session-depth' import { withDevinSessionsDbScan } from './session-scanner-devin-db' +import { withOpenCodeSqliteScanScope } from './session-scanner-opencode-sqlite-scan-scope' const SESSION_PARSE_CONCURRENCY = 8 const SESSION_PARSE_CANDIDATE_MULTIPLIER = 2 @@ -61,6 +62,10 @@ const SESSION_PARSE_CANDIDATE_MULTIPLIER = 2 export async function scanAiVaultSessions( options: AiVaultScanOptions = {} ): Promise { + return withOpenCodeSqliteScanScope(() => scanAiVaultSessionStores(options)) +} + +async function scanAiVaultSessionStores(options: AiVaultScanOptions): Promise { // The span makes scan cost visible in the local trace file: STA-1278-style // "one core pegged" reports need to show whether transcript scanning is the // subsystem burning CPU, and how much of each scan the cache absorbed. diff --git a/src/main/codex/codex-collab-call-delegation.test.ts b/src/main/codex/codex-collab-call-delegation.test.ts index 96131248ddc..dbc116802e4 100644 --- a/src/main/codex/codex-collab-call-delegation.test.ts +++ b/src/main/codex/codex-collab-call-delegation.test.ts @@ -28,7 +28,9 @@ import { THREAD_ID } from './codex-structured-session-adapter-fixture' /** The delegation the parent's newest tool run reads as, the row a running chat's frontier judges. */ async function newestRunDelegation(frames: Frame[]): Promise { const { conversation } = projectNativeChatTranscript( - projectStructuredAgentSessionMessages(await publishedRows(frames), [], []) + projectStructuredAgentSessionMessages(await publishedRows(frames), [], [], { + rejectedInPlace: true + }) ) const runs = conversation.filter((message: NativeChatMessage) => message.blocks.some((block) => block.type === 'tool-call') diff --git a/src/main/codex/codex-collab-call-rows-on-positional-clients.test.ts b/src/main/codex/codex-collab-call-rows-on-positional-clients.test.ts index 795716ed553..ae7770b5beb 100644 --- a/src/main/codex/codex-collab-call-rows-on-positional-clients.test.ts +++ b/src/main/codex/codex-collab-call-rows-on-positional-clients.test.ts @@ -47,7 +47,9 @@ function positionalRuns(rows: AgentJournalRenderItem[]): { }) }) const transcript = projectNativeChatTranscript( - projectStructuredAgentSessionMessages(rows, [], []).map(withoutCallIds) + projectStructuredAgentSessionMessages(rows, [], [], { rejectedInPlace: true }).map( + withoutCallIds + ) ) // The conversation's runs, then each helper's section's. const runs = [ diff --git a/src/main/codex/codex-structured-question-order.test.ts b/src/main/codex/codex-structured-question-order.test.ts index 8ad4b75bd38..78c197644e2 100644 --- a/src/main/codex/codex-structured-question-order.test.ts +++ b/src/main/codex/codex-structured-question-order.test.ts @@ -206,6 +206,7 @@ function drawnPromptRows(): string[][] { client.items, [], client.submissions, + { rejectedInPlace: true }, projectStructuredQuestionMessages ) ) @@ -243,9 +244,9 @@ describe('a Codex ask with several questions', () => { ]) // Mobile draws the shared projection in journal order, one row per question. expect( - projectStructuredAgentSessionMessages(client.items, [], client.submissions).map( - ({ blocks }) => (blocks[0]?.type === 'text' ? blocks[0].text.split('\n')[0] : null) - ) + projectStructuredAgentSessionMessages(client.items, [], client.submissions, { + rejectedInPlace: false + }).map(({ blocks }) => (blocks[0]?.type === 'text' ? blocks[0].text.split('\n')[0] : null)) ).toEqual(ASKED.map(({ question }) => question)) }) diff --git a/src/main/native-chat/agent-model-catalog/agent-model-catalog-service.test.ts b/src/main/native-chat/agent-model-catalog/agent-model-catalog-service.test.ts index d51a98b3d8b..02025485151 100644 --- a/src/main/native-chat/agent-model-catalog/agent-model-catalog-service.test.ts +++ b/src/main/native-chat/agent-model-catalog/agent-model-catalog-service.test.ts @@ -5,7 +5,11 @@ import { agentModelCatalogFingerprintForRecord } from './agent-model-catalog-fingerprint' import { createAgentModelCatalogService } from './agent-model-catalog-service' -import { AgentModelCatalogStore, type AgentModelCatalogSuccess } from './agent-model-catalog-store' +import { + AGENT_MODEL_CATALOG_FRESH_MS, + AgentModelCatalogStore, + type AgentModelCatalogSuccess +} from './agent-model-catalog-store' function record(accountHomePath: string): AgentSessionRecord { // oxlint-disable-next-line typescript/consistent-type-assertions -- SAFETY: the service reads only provider, accountHome and location; the rest of the record is irrelevant here. @@ -55,7 +59,8 @@ describe('agent model catalog service', () => { probes: { codex: probe } }) expect(await service.read({ agent: 'codex', sessionId: 'session-1' })).toEqual({ - origin: 'unknown' + origin: 'unknown', + listingInProgress: true }) // A second read while the probe is in flight must not start another, and a // record-scoped read probes the RECORD's pinned home, not the selection. @@ -81,7 +86,10 @@ describe('agent model catalog service', () => { probes: { codex: probe } }) // The record-less read follows the CURRENT selection: unknown, never gpt-old. - expect(await service.read({ agent: 'codex' })).toEqual({ origin: 'unknown' }) + expect(await service.read({ agent: 'codex' })).toEqual({ + origin: 'unknown', + listingInProgress: true + }) expect(probe).toHaveBeenCalledWith('/homes/new') await vi.waitFor(async () => { const result = await service.read({ agent: 'codex' }) @@ -139,7 +147,8 @@ describe('agent model catalog service', () => { probes: { codex: probe } }) expect(await service.read({ agent: 'codex', sessionId: 'session-1' })).toEqual({ - origin: 'unknown' + origin: 'unknown', + listingInProgress: true }) await vi.waitFor(() => expect(probe).toHaveBeenCalledTimes(1)) // Still a clean unknown — and the failure TTL suppresses a probe storm. @@ -164,6 +173,94 @@ describe('agent model catalog service', () => { expect(probe).not.toHaveBeenCalled() }) + describe('a read that waits for the first listing', () => { + function deferredListing() { + let resolve!: (success: AgentModelCatalogSuccess) => void + let reject!: (error: Error) => void + const promise = new Promise((res, rej) => { + resolve = res + reject = rej + }) + return { promise, resolve, reject } + } + + function coldService(probe: (home: string) => Promise) { + const store = new AgentModelCatalogStore() + const service = createAgentModelCatalogService({ + store, + getRecord: () => undefined, + resolveAccountHome: async () => CODEX_HOME('/homes/selected'), + probes: { codex: probe } + }) + return { store, service } + } + + it('joins the listing the first read started and answers with it', async () => { + const pending = deferredListing() + const probe = vi.fn(() => pending.promise) + const { service } = coldService(probe) + expect(await service.read({ agent: 'codex' })).toEqual({ + origin: 'unknown', + listingInProgress: true + }) + const waited = service.read({ agent: 'codex', waitForListing: true }) + pending.resolve(listing('gpt-listed')) + const result = await waited + expect(result.origin === 'unknown' ? null : result.models[0]!.id).toBe('gpt-listed') + expect(probe).toHaveBeenCalledTimes(1) + }) + + it('answers a plain unknown when the listing fails', async () => { + const pending = deferredListing() + const { service } = coldService(() => pending.promise) + const waited = service.read({ agent: 'codex', waitForListing: true }) + pending.reject(new Error('spawn failed')) + expect(await waited).toEqual({ origin: 'unknown' }) + }) + + it('does not wait or report a listing while a failure is inside its TTL', async () => { + const probe = vi.fn(async (): Promise => { + throw new Error('spawn failed') + }) + const { store, service } = coldService(probe) + store.recordFailure(selectedHomeFingerprint('/homes/selected'), 'spawn failed') + expect(await service.read({ agent: 'codex', waitForListing: true })).toEqual({ + origin: 'unknown' + }) + expect(await service.read({ agent: 'codex' })).toEqual({ origin: 'unknown' }) + expect(probe).not.toHaveBeenCalled() + }) + + it('reports no listing where the host has no lister for the account', async () => { + const store = new AgentModelCatalogStore() + const service = createAgentModelCatalogService({ + store, + getRecord: () => undefined, + resolveAccountHome: async () => CODEX_HOME('/homes/selected') + }) + expect(await service.read({ agent: 'codex', waitForListing: true })).toEqual({ + origin: 'unknown' + }) + }) + + it('serves an aged entry at once and refreshes it behind the answer', async () => { + let now = 0 + const store = new AgentModelCatalogStore({ now: () => now }) + store.recordSuccess(selectedHomeFingerprint('/homes/selected'), 'codex', listing('gpt-old')) + now = AGENT_MODEL_CATALOG_FRESH_MS + const probe = vi.fn(() => new Promise(() => {})) + const service = createAgentModelCatalogService({ + store, + getRecord: () => undefined, + resolveAccountHome: async () => CODEX_HOME('/homes/selected'), + probes: { codex: probe } + }) + const result = await service.read({ agent: 'codex', waitForListing: true }) + expect(result.origin === 'unknown' ? null : result.models[0]!.id).toBe('gpt-old') + expect(probe).toHaveBeenCalledTimes(1) + }) + }) + describe('a read for the workspace a new chat runs in', () => { function serviceWith(mayOverride: boolean) { const store = new AgentModelCatalogStore() diff --git a/src/main/native-chat/agent-model-catalog/agent-model-catalog-service.ts b/src/main/native-chat/agent-model-catalog/agent-model-catalog-service.ts index 8c794331a2c..638b0f54272 100644 --- a/src/main/native-chat/agent-model-catalog/agent-model-catalog-service.ts +++ b/src/main/native-chat/agent-model-catalog/agent-model-catalog-service.ts @@ -35,6 +35,8 @@ export type AgentModelCatalogService = { sessionId?: string /** Where a new chat would run; null when one was named but is not a local directory. */ workspacePath?: string | null + /** With no entry yet, answer from the listing this read starts or joins instead of `unknown`. */ + waitForListing?: boolean }) => Promise } @@ -79,8 +81,9 @@ async function workspaceKeepsListedDefault( * launch); without one, the key is the account a launch would pin right now — * never "whichever account listed last". `unknown` tells the client to keep * its static seed, and a missing or aged entry kicks one joined background - * probe so the next read is warm. Failures are the store's 30s TTL, never an - * answer — a picker is a user surface and must not block. + * probe so the next read is warm. With no entry, the answer says that listing + * is running, and only a read that asks waits for it. Failures are the store's + * 30s TTL, never an answer: inside it a read answers `unknown` at once. */ export function createAgentModelCatalogService( deps: AgentModelCatalogServiceDeps @@ -110,14 +113,27 @@ export function createAgentModelCatalogService( }) accountHomePath = resolved.path } - const entry = deps.store.get(fingerprint) + let entry = deps.store.get(fingerprint) const probe = deps.probes?.[params.agent] - if (probe && accountHomePath && deps.store.shouldRefresh(fingerprint)) { - const home = accountHomePath - void deps.store.refresh(fingerprint, params.agent, () => probe(home)) - } + const home = accountHomePath + // Without an entry, join a running listing too: that is the one a waiting read answers from. + const listing = + probe && + home && + (entry ? deps.store.shouldRefresh(fingerprint) : !deps.store.hasActiveFailure(fingerprint)) + ? deps.store.refresh(fingerprint, params.agent, () => probe(home)) + : null if (!entry) { - return { origin: 'unknown' } + if (!listing) { + return { origin: 'unknown' } + } + if (!params.waitForListing) { + return { origin: 'unknown', listingInProgress: true } + } + entry = await listing + if (!entry) { + return { origin: 'unknown' } + } } return resultFromEntry( entry, diff --git a/src/main/native-chat/agent-session-journal/journal-dispatch-reducer.ts b/src/main/native-chat/agent-session-journal/journal-dispatch-reducer.ts index 89722959799..c7c2a671654 100644 --- a/src/main/native-chat/agent-session-journal/journal-dispatch-reducer.ts +++ b/src/main/native-chat/agent-session-journal/journal-dispatch-reducer.ts @@ -8,7 +8,11 @@ import { import { agentJournalSubmissionKey } from '../../../shared/agent-session-journal-item-key' import { journalDispatchRowApplies } from './journal-dispatch-settlement' import type { JournalReducerState } from './journal-reducer' -import { notePersonTurnAccepted, placeHandedOverMessage } from './journal-submission-fold' +import { + notePersonTurnAccepted, + placeHandedOverMessage, + placeRejectedMessage +} from './journal-submission-fold' import type { JournalRow } from './journal-row-schema' export function applyJournalDispatchRow( @@ -34,6 +38,8 @@ export function applyJournalDispatchRow( if (row.state === 'pending') { submission.handedOverAt = row.ts placeHandedOverMessage(state, submission, row) + } else if (row.state === 'rejected') { + placeRejectedMessage(state, submission, row) } if (row.recovered) { submission.recovered = row.recovered diff --git a/src/main/native-chat/agent-session-journal/journal-lifecycle-batch-appender.ts b/src/main/native-chat/agent-session-journal/journal-lifecycle-batch-appender.ts index c05d22a8745..68ba188e402 100644 --- a/src/main/native-chat/agent-session-journal/journal-lifecycle-batch-appender.ts +++ b/src/main/native-chat/agent-session-journal/journal-lifecycle-batch-appender.ts @@ -3,6 +3,7 @@ import type { JournalReducerState } from './journal-reducer' import { journalLifecycleBatchRowBuilder } from './journal-row-builders' import type { JournalLifecycleBatchInput } from './journal-store-contracts' import type { JournalRow } from './journal-row-schema' +import { journalQueuedRejectionRowBuilders } from './journal-pending-submission-recovery' const SETTLEMENT_ALREADY_APPLIED = new Error('journal_settlement_already_applied') @@ -12,10 +13,38 @@ export class JournalLifecycleBatchAppender { state: () => JournalReducerState cursor: () => AgentJournalCursor enqueue: (build: (seq: number, ts: number) => JournalRow) => Promise + enqueueRows: ( + plan: () => readonly ((seq: number, ts: number) => JournalRow)[] + ) => Promise } ) {} append(input: JournalLifecycleBatchInput): Promise { + const { rejectsQueued } = input + if (rejectsQueued) { + // Planned on the lane: the sends queued then, and this batch unless it already landed. With + // none left (a Stop withdrew them first) it failed no one, so nothing is written. + return this.deps + .enqueueRows(() => { + const rejections = journalQueuedRejectionRowBuilders( + this.deps.state, + input.fence, + rejectsQueued + ) + return rejections.length === 0 || this.wasApplied(input.settlementId) + ? rejections + : [ + ...rejections, + journalLifecycleBatchRowBuilder( + this.deps.state, + input.settlementId, + input.mutations, + input + ) + ] + }) + .then(() => this.deps.cursor()) + } if (this.wasApplied(input.settlementId)) { return Promise.resolve(this.deps.cursor()) } diff --git a/src/main/native-chat/agent-session-journal/journal-pending-submission-recovery.ts b/src/main/native-chat/agent-session-journal/journal-pending-submission-recovery.ts index 964ea42f931..b52e3648b72 100644 --- a/src/main/native-chat/agent-session-journal/journal-pending-submission-recovery.ts +++ b/src/main/native-chat/agent-session-journal/journal-pending-submission-recovery.ts @@ -2,6 +2,9 @@ import type { AgentJournalDispatchRejection } from '../../../shared/agent-sessio import type { AgentJournalSubmission } from '../../../shared/agent-session-journal-types' import { isQueuedAgentJournalSubmission } from '../../../shared/agent-session-queued-submission' import { DISPATCH_DOUBT_HOST_RESTARTED } from './journal-dispatch-doubt-reasons' +import type { JournalReducerState } from './journal-reducer' +import { journalDispatchRowBuilder } from './journal-row-builders' +import type { JournalRow } from './journal-row-schema' import type { AgentSessionJournal } from './journal-store' /** Settles every submission a process fact left unanswerable. Doubt is never @@ -88,3 +91,21 @@ export async function rejectJournalQueuedSubmissions( ) return queued.map((entry) => entry.clientMessageId) } + +/** Rows rejecting every submission still queued, read from `state` when called: for an append + * that must carry them with what follows, in one transaction. */ +export function journalQueuedRejectionRowBuilders( + state: () => JournalReducerState, + fence: number, + rejection: AgentJournalDispatchRejection +): ((seq: number, ts: number) => JournalRow)[] { + return [...state().submissions.values()].filter(isQueuedAgentJournalSubmission).map((entry) => + journalDispatchRowBuilder(state, { + clientMessageId: entry.clientMessageId, + state: 'rejected', + ...rejection, + fence, + recovered: true + }) + ) +} diff --git a/src/main/native-chat/agent-session-journal/journal-queued-rejection-batch.test.ts b/src/main/native-chat/agent-session-journal/journal-queued-rejection-batch.test.ts new file mode 100644 index 00000000000..d218a356fae --- /dev/null +++ b/src/main/native-chat/agent-session-journal/journal-queued-rejection-batch.test.ts @@ -0,0 +1,129 @@ +// A failed start's row and the queued messages it failed land in ONE append, the messages first: +// no reader meets one without the other, and the messages sit above the row that says why. + +import { mkdtemp, rm } from 'node:fs/promises' +import { tmpdir } from 'node:os' +import { join } from 'node:path' +import { afterEach, beforeEach, expect, it } from 'vitest' +import { agentSessionFailureFact } from '../../../shared/agent-session-failure' +import { agentSessionFailureWords } from '../../../shared/agent-session-failure-words' +import { agentJournalSubmissionKey } from '../../../shared/agent-session-journal-item-key' +import { + AGENT_JOURNAL_THREAD_SCOPE, + type AgentSessionJournalIdentity +} from '../../../shared/agent-session-journal-types' +import type { AgentSessionJournal } from './journal-store' +import { + closeTestJournalHostDatabases, + createTrackedJournalOpener +} from './journal-host-database-test-support' + +const IDENTITY: AgentSessionJournalIdentity = { + sessionId: 'session-start', + workspaceId: 'ws-1', + hostId: 'host-1', + agent: 'claude', + providerHandle: { kind: 'claude', sessionId: 'native-1', leafUuid: null } +} + +const START_FAILED = agentSessionFailureWords(agentSessionFailureFact('providerStartFailed'), { + surface: 'rejection' +}) +const ERROR_ROW = { provider: 'orca', clientMessageId: 'start-failure:gen-1' } as const + +let root: string +let clock = 1_000 +const journals = createTrackedJournalOpener() + +beforeEach(async () => { + root = await mkdtemp(join(tmpdir(), 'orca-queued-rejection-batch-')) +}) + +afterEach(async () => { + await closeTestJournalHostDatabases() + await rm(root, { recursive: true, force: true }) +}) + +async function openWithQueued(...ids: string[]): Promise { + const journal = await journals.open({ + identity: IDENTITY, + stateDirectory: root, + now: () => (clock += 1), + mintEpoch: () => 'epoch-1' + }) + for (const id of ids) { + await journal.appendSubmission({ + clientMessageId: id, + payloadFingerprint: `fp-${id}`, + body: { kind: 'message', role: 'user', blocks: [{ type: 'text', text: id }] }, + fence: 0, + handoverRecorded: true + }) + } + return journal +} + +function startFailureBatch(mutations = 1) { + return { + settlementId: 'start-failure:gen-1', + fence: 0, + recovered: true as const, + mutations: Array.from({ length: mutations }, () => ({ + kind: 'item' as const, + identity: ERROR_ROW, + body: { kind: 'status' as const, tone: 'error' as const, text: 'Claude did not start.' }, + turnScope: AGENT_JOURNAL_THREAD_SCOPE + })), + rejectsQueued: START_FAILED + } +} + +it('writes the rejections first and the row after them, and draws the messages above it', async () => { + const journal = await openWithQueued('first', 'second') + const before = journal.cursor().sequence + + await journal.appendLifecycleBatch(startFailureBatch()) + + expect(journal.cursor().sequence).toBe(before + 3) + expect(journal.submissions().map((entry) => entry.dispatchState)).toEqual([ + 'rejected', + 'rejected' + ]) + const order = journal.snapshot().items.map((item) => item.itemId) + expect(order).toEqual([ + agentJournalSubmissionKey('first'), + agentJournalSubmissionKey('second'), + 'orca:start-failure%3Agen-1' + ]) +}) + +it('writes neither when the row cannot be written', async () => { + const journal = await openWithQueued('first') + const before = journal.cursor().sequence + + // An empty batch breaks the row's bound, so the transaction rolls back as a whole. + await expect(journal.appendLifecycleBatch(startFailureBatch(0))).rejects.toThrow( + 'journal_lifecycle_batch_mutation_bound_exceeded' + ) + + expect(journal.cursor().sequence).toBe(before) + expect(journal.submissions().map((entry) => entry.dispatchState)).toEqual(['pending']) +}) + +// A Stop that reaches the lane first takes the message back; the failed start then failed no one. +it('writes nothing when a Stop withdrew every queued message first', async () => { + const journal = await openWithQueued('first') + const withdrawal = agentSessionFailureWords(agentSessionFailureFact('cancelled'), { + surface: 'rejection' + }) + + await Promise.all([ + journal.rejectQueuedSubmissions(0, withdrawal), + journal.appendLifecycleBatch(startFailureBatch()) + ]) + + expect(journal.submissions()[0]?.rejection).toEqual({ kind: 'cancelled' }) + expect(journal.snapshot().items.map((item) => item.itemId)).toEqual([ + agentJournalSubmissionKey('first') + ]) +}) diff --git a/src/main/native-chat/agent-session-journal/journal-reducer.test.ts b/src/main/native-chat/agent-session-journal/journal-reducer.test.ts index 45681f32871..60d068ebfd8 100644 --- a/src/main/native-chat/agent-session-journal/journal-reducer.test.ts +++ b/src/main/native-chat/agent-session-journal/journal-reducer.test.ts @@ -551,7 +551,7 @@ describe('submission and dispatch state machine', () => { expect(state.receipts.get('cm_1')).toBeTruthy() }) - it('keeps a refused write rejected and leaves its bubble where it was', () => { + it('keeps a refused write rejected, at its rejection, whatever comes after', () => { const state = fold([ submission, { @@ -579,7 +579,8 @@ describe('submission and dispatch state machine', () => { submittedAt: submission.ts, reason: 'provider_write_failed: closed before enqueue' }) - expect(renderJournalState(state).items[0]?.sequence).toBe(submission.seq) + // It sits where it was rejected; the late `pending` moves nothing. + expect(renderJournalState(state).items[0]?.sequence).toBe(2) }) it('ignores a dispatch for a submission this epoch never saw', () => { diff --git a/src/main/native-chat/agent-session-journal/journal-reducer.ts b/src/main/native-chat/agent-session-journal/journal-reducer.ts index a7d6f150e70..e2d1b1ef3ca 100644 --- a/src/main/native-chat/agent-session-journal/journal-reducer.ts +++ b/src/main/native-chat/agent-session-journal/journal-reducer.ts @@ -6,9 +6,9 @@ // dropped rather than resurrecting stale content, and ordering is by the // position (sequence, then place in the row) of the write that CREATED an item // (a later revision updates the body, it does not move the bubble) — except a -// queued message, which sits where its handover put it. Producer linkage is -// likewise the creating write's: a revision naming no producer keeps it, one -// naming any replaces it. +// queued message, which sits where its handover put it, and a rejected one, which +// sits where it was rejected. Producer linkage is likewise the creating write's: a +// revision naming no producer keeps it, one naming any replaces it. import type { AgentJournalAcceptanceReceipt, diff --git a/src/main/native-chat/agent-session-journal/journal-rejected-message-placement.test.ts b/src/main/native-chat/agent-session-journal/journal-rejected-message-placement.test.ts new file mode 100644 index 00000000000..ce5b407d3e8 --- /dev/null +++ b/src/main/native-chat/agent-session-journal/journal-rejected-message-placement.test.ts @@ -0,0 +1,158 @@ +// A rejected message sits where it was rejected, in no turn: queued, handed over or sent directly. +// One in doubt stays where it was: it may have reached the agent. + +import { describe, expect, it } from 'vitest' +import { agentJournalSubmissionKey } from '../../../shared/agent-session-journal-item-key' +import { + AGENT_JOURNAL_THREAD_SCOPE, + type AgentJournalTurnScope +} from '../../../shared/agent-session-journal-types' +import { DISPATCH_REJECTED_HOST_RESTARTED } from '../../../shared/structured-agent-session-dispatch-rejection' +import { projectStructuredAgentSessionMessages } from '../../../shared/structured-agent-session-message-projection' +import { projectJournalBatch } from '../agent-session-wire/agent-session-journal-batch' +import { DISPATCH_DOUBT_HOST_RESTARTED } from './journal-dispatch-doubt-reasons' +import { applyJournalRow, createJournalReducerState, renderJournalState } from './journal-reducer' +import { buildJournalSubmissionRow, journalRowBase } from './journal-row-builders' +import type { JournalRow } from './journal-row-schema' + +function journal() { + const state = createJournalReducerState('session-1', 'epoch-1') + let seq = 0 + const push = (next: JournalRow): JournalRow => { + applyJournalRow(state, next) + return next + } + return { + state, + submission(clientMessageId: string, handoverRecorded = true) { + seq += 1 + return push( + buildJournalSubmissionRow({ + state, + clientMessageId, + payloadFingerprint: `fp-${clientMessageId}`, + providerHandle: { kind: 'codex', threadId: 'thread-1' }, + body: { + kind: 'message', + role: 'user', + blocks: [{ type: 'text', text: clientMessageId }] + }, + seq, + fence: 1, + ts: 1_000 + seq, + ...(handoverRecorded ? { handoverRecorded: true } : {}) + }) + ) + }, + dispatch( + clientMessageId: string, + state_: 'pending' | 'rejected' | 'unknown', + turnScope: AgentJournalTurnScope = AGENT_JOURNAL_THREAD_SCOPE + ) { + seq += 1 + return push({ + kind: 'dispatch', + clientMessageId, + state: state_, + providerItemId: null, + reason: + state_ === 'rejected' + ? DISPATCH_REJECTED_HOST_RESTARTED + : state_ === 'unknown' + ? DISPATCH_DOUBT_HOST_RESTARTED + : null, + ...(state_ === 'rejected' ? { rejection: { kind: 'hostRestarted' } } : {}), + ...journalRowBase(state.epoch, seq, 1, 1_000 + seq), + turnScope + }) + }, + /** Other work between the send and its settlement, as a turn running ahead of it would write. */ + advance(rows: number) { + seq += rows + state.lastSequence = seq + } + } +} + +function placed(state: ReturnType['state'], clientMessageId: string) { + const item = state.items.get(agentJournalSubmissionKey(clientMessageId)) + return item && { sequence: item.sequence, observedAt: item.observedAt, scope: item.turnScope } +} + +const IN_TURN: AgentJournalTurnScope = { kind: 'turn', turnItemId: 'orca:turn-1' } + +describe('a rejected message', () => { + it('sits at the rejection, in no turn, when it was queued', () => { + const { state, submission, dispatch, advance } = journal() + submission('waiting') + advance(300) + dispatch('waiting', 'rejected') + + expect(placed(state, 'waiting')).toEqual({ + sequence: 302, + observedAt: 1_302, + scope: AGENT_JOURNAL_THREAD_SCOPE + }) + }) + + it('reaches a subscriber at the tail, so the newest page holds it', () => { + const { state, submission, dispatch, advance } = journal() + submission('waiting') + advance(300) + const rejection = dispatch('waiting', 'rejected') + + const projected = projectJournalBatch({ + rows: [rejection], + snapshot: renderJournalState(state), + afterSequence: 301 + }) + expect(projected.ok && projected.batch.items.map((item) => item.sequence)).toEqual([302]) + expect(renderJournalState(state).items.at(-1)?.itemId).toBe( + agentJournalSubmissionKey('waiting') + ) + }) + + it('sits at the rejection, out of the turn it was handed into, when it was a steer', () => { + const { state, submission, dispatch, advance } = journal() + submission('steer') + advance(10) + dispatch('steer', 'pending', IN_TURN) + expect(placed(state, 'steer')?.scope).toEqual(IN_TURN) + advance(10) + dispatch('steer', 'rejected') + + expect(placed(state, 'steer')).toEqual({ + sequence: 23, + observedAt: 1_023, + scope: AGENT_JOURNAL_THREAD_SCOPE + }) + }) + + it('sits at the rejection when it was sent directly', () => { + const { state, submission, dispatch, advance } = journal() + submission('direct', false) + advance(10) + dispatch('direct', 'rejected') + + expect(placed(state, 'direct')?.sequence).toBe(12) + }) +}) + +describe('a message in doubt', () => { + it('stays where it was handed over, a plain bubble in its turn: it may have reached the agent', () => { + const { state, submission, dispatch, advance } = journal() + submission('steer') + advance(10) + dispatch('steer', 'pending', IN_TURN) + advance(10) + dispatch('steer', 'unknown') + + expect(placed(state, 'steer')).toEqual({ sequence: 12, observedAt: 1_012, scope: IN_TURN }) + const { items, submissions } = renderJournalState(state) + const drawn = projectStructuredAgentSessionMessages(items, [], submissions, { + rejectedInPlace: true + }).find((message) => message.id === agentJournalSubmissionKey('steer')) + expect(drawn).toBeDefined() + expect(drawn?.unsent).toBeUndefined() + }) +}) diff --git a/src/main/native-chat/agent-session-journal/journal-row-writer.test.ts b/src/main/native-chat/agent-session-journal/journal-row-writer.test.ts index c0824a2e4ca..3671ce46cfc 100644 --- a/src/main/native-chat/agent-session-journal/journal-row-writer.test.ts +++ b/src/main/native-chat/agent-session-journal/journal-row-writer.test.ts @@ -93,4 +93,17 @@ describe('journal row writer', () => { kind: 'item' }) }) + + it('writes several rows as one: a failure on the last leaves none', async () => { + const { writer, committedRows } = writerHarness() + // Sequence 2 is taken, so the second of the two rows violates the primary key. + insertTestJournalRow(database.db, SESSION_ID, row(2, 1)) + + await expect(writer.enqueueRows(() => [row, row])).rejects.toThrow() + + expect(committedRows).toHaveLength(0) + expect(readTestJournalRows(database.db, SESSION_ID, EPOCH).map((stored) => stored.seq)).toEqual( + [2] + ) + }) }) diff --git a/src/main/native-chat/agent-session-journal/journal-row-writer.ts b/src/main/native-chat/agent-session-journal/journal-row-writer.ts index ca274120a7b..dde1a66741d 100644 --- a/src/main/native-chat/agent-session-journal/journal-row-writer.ts +++ b/src/main/native-chat/agent-session-journal/journal-row-writer.ts @@ -71,6 +71,40 @@ export class JournalRowWriter { }) } + /** Several rows in ONE transaction, in order, planned once the lane is this append's: none is + * durable unless all are, so no reader ever meets some without the rest. */ + enqueueRows( + plan: () => readonly ((seq: number, ts: number) => JournalRow)[] + ): Promise { + return this.deps.serialize(() => { + assertJournalWritable(this.deps.readOnly(), this.deps.sessionId) + const first = this.deps.nextSequence() + const ts = this.deps.now() + const rows = plan().map((build, index) => build(first + index, ts)) + if (rows.length === 0) { + return rows + } + for (const row of rows) { + assertJournalFence(row.fence, this.deps.highestFence()) + } + try { + this.deps.database().transaction((db) => { + for (const row of rows) { + insertJournalRow(db, this.deps.sessionId, row) + this.runBookkeeping(db, row) + } + }) + } catch (error) { + this.deps.rolledBack?.() + throw error + } + for (const row of rows) { + this.deps.commit(row) + } + return rows + }) + } + /** Assign the next sequence, make the row durable, and fold it through the SAME reducer * replay uses — all inside one serialized step — answering where the row landed. */ append( diff --git a/src/main/native-chat/agent-session-journal/journal-store-collaborators.ts b/src/main/native-chat/agent-session-journal/journal-store-collaborators.ts index 65ea665879b..ba879696469 100644 --- a/src/main/native-chat/agent-session-journal/journal-store-collaborators.ts +++ b/src/main/native-chat/agent-session-journal/journal-store-collaborators.ts @@ -92,6 +92,20 @@ export function createJournalStoreCollaborators(host: JournalStoreHost): Journal wroteBeforeOpen: (sequence) => host.journal().wroteBeforeOpen(sequence), committed: host.notifyCommitted }) + const rowWriter = new JournalRowWriter({ + sessionId: host.identity.sessionId, + now: host.now, + serialize: host.serialize, + database: host.database, + readOnly: host.readOnly, + highestFence: () => host.state().highestFence, + nextSequence: () => host.state().lastSequence + 1, + commit: host.commit, + // Every rejection is a dispatch row through this one writer; the draft + // returned-transition rides it so no path can bypass the hook. + inTransaction: (db, row) => queuedMessages.onRowInTransaction(db, row), + rolledBack: () => queuedMessages.invalidate() + }) return { epochController, queuedMessages, @@ -102,20 +116,7 @@ export function createJournalStoreCollaborators(host: JournalStoreHost): Journal restoreJournalStore(host, { epochController }).then(() => queuedMessages.repairAndPruneAtOpen() ), - rowWriter: new JournalRowWriter({ - sessionId: host.identity.sessionId, - now: host.now, - serialize: host.serialize, - database: host.database, - readOnly: host.readOnly, - highestFence: () => host.state().highestFence, - nextSequence: () => host.state().lastSequence + 1, - commit: host.commit, - // Every rejection is a dispatch row through this one writer; the draft - // returned-transition rides it so no path can bypass the hook. - inTransaction: (db, row) => queuedMessages.onRowInTransaction(db, row), - rolledBack: () => queuedMessages.invalidate() - }), + rowWriter, itemAppender: new JournalItemAppender({ state: host.state, enqueue: host.enqueue @@ -123,7 +124,8 @@ export function createJournalStoreCollaborators(host: JournalStoreHost): Journal lifecycleBatchAppender: new JournalLifecycleBatchAppender({ state: host.state, cursor: host.cursor, - enqueue: host.enqueue + enqueue: host.enqueue, + enqueueRows: (plan) => rowWriter.enqueueRows(plan) }) } } diff --git a/src/main/native-chat/agent-session-journal/journal-store-contracts.ts b/src/main/native-chat/agent-session-journal/journal-store-contracts.ts index 85e6b99f902..cb86a7a8e5d 100644 --- a/src/main/native-chat/agent-session-journal/journal-store-contracts.ts +++ b/src/main/native-chat/agent-session-journal/journal-store-contracts.ts @@ -70,6 +70,10 @@ export type JournalLifecycleBatchInput = { mutations: readonly JournalLifecycleMutationInput[] fence: number recovered?: true + /** Rejects the sends still queued with this first, in the same append: a failed start's row + * follows the messages it failed, and no reader meets one without the other. With none still + * queued, the batch is not written either. */ + rejectsQueued?: AgentJournalDispatchRejection } export type JournalSubmissionInput = { diff --git a/src/main/native-chat/agent-session-journal/journal-submission-fold.ts b/src/main/native-chat/agent-session-journal/journal-submission-fold.ts index 2c9c6cfcd04..80e4e2c93e9 100644 --- a/src/main/native-chat/agent-session-journal/journal-submission-fold.ts +++ b/src/main/native-chat/agent-session-journal/journal-submission-fold.ts @@ -69,6 +69,28 @@ export function placeHandedOverMessage( }) } +/** A rejected message — queued, handed over, or sent directly — joins the conversation where it was + * rejected, in no turn: what happened before the rejection happened before it, and the newest page + * holds a recent one. Only a rejection: one in doubt may have reached the agent, so it stays. */ +export function placeRejectedMessage( + state: JournalReducerState, + submission: AgentJournalSubmission, + row: Extract +): void { + const itemId = agentJournalSubmissionKey(submission.clientMessageId) + const item = state.items.get(itemId) + if (row.state !== 'rejected' || !item) { + return + } + const { sequenceIndex: _placed, ...rest } = item + state.items.set(itemId, { + ...rest, + sequence: row.seq, + observedAt: row.ts, + turnScope: AGENT_JOURNAL_THREAD_SCOPE + }) +} + export function acceptSubmissionFromProviderItem( state: JournalReducerState, providerItemId: string, diff --git a/src/main/native-chat/agent-session-journal/journal-turn-scope.test.ts b/src/main/native-chat/agent-session-journal/journal-turn-scope.test.ts index 2b8e4b80a76..4674f6de929 100644 --- a/src/main/native-chat/agent-session-journal/journal-turn-scope.test.ts +++ b/src/main/native-chat/agent-session-journal/journal-turn-scope.test.ts @@ -163,7 +163,9 @@ describe('stated turn scope', () => { const expected = [agentJournalItemKey(row('result')), agentJournalSubmissionKey('held')] const onPhone = () => { const { items, submissions } = renderJournalState(state) - return projectStructuredAgentSessionMessages(items, [], submissions) + return projectStructuredAgentSessionMessages(items, [], submissions, { + rejectedInPlace: false + }) } // Still waiting: drawn after everything the agent did, the command's result included. expect(drawn(onPhone())).toEqual(expected) diff --git a/src/main/native-chat/agent-session-wire/structured-agent-session-conversation-close.test.ts b/src/main/native-chat/agent-session-wire/structured-agent-session-conversation-close.test.ts index c084d2b2e07..2dac4cacb7d 100644 --- a/src/main/native-chat/agent-session-wire/structured-agent-session-conversation-close.test.ts +++ b/src/main/native-chat/agent-session-wire/structured-agent-session-conversation-close.test.ts @@ -241,16 +241,16 @@ describe('a start that never finishes (P2-15)', () => { providerChildPhase: 'starting' as const })) Object.assign(rig.host.deps.adapter, { awaitStarted: () => started.promise }) - const reject = AgentSessionJournal.prototype.rejectQueuedSubmissions - vi.spyOn(AgentSessionJournal.prototype, 'rejectQueuedSubmissions').mockImplementation(function ( + // The loop rejects the queued messages in the same append as its row. + const append = AgentSessionJournal.prototype.appendLifecycleBatch + vi.spyOn(AgentSessionJournal.prototype, 'appendLifecycleBatch').mockImplementation(function ( this: AgentSessionJournal, ...args ) { - // Not the open's sweep of an earlier process's leftovers. - if (args[1].rejection.kind !== 'hostRestarted') { - order.push(`rejected: ${args[1].reason}`) + if (args[0].rejectsQueued) { + order.push(`rejected: ${args[0].rejectsQueued.reason}`) } - return reject.apply(this, args) + return append.apply(this, args) }) const reader = collectSubscriber() const attached = await rig.host.attach(CALLER, hostTestAttachParams(null)) diff --git a/src/main/native-chat/agent-session-wire/structured-agent-session-host-test-data.ts b/src/main/native-chat/agent-session-wire/structured-agent-session-host-test-data.ts index 959a02919a2..a96388ac4ca 100644 --- a/src/main/native-chat/agent-session-wire/structured-agent-session-host-test-data.ts +++ b/src/main/native-chat/agent-session-wire/structured-agent-session-host-test-data.ts @@ -84,6 +84,8 @@ export function hostTestDrawnRowIds( state: 'dispatching' as const })) return projectNativeChatTranscriptMessages( - projectStructuredAgentSessionMessages(snapshot.items, outbox, snapshot.submissions) + projectStructuredAgentSessionMessages(snapshot.items, outbox, snapshot.submissions, { + rejectedInPlace: true + }) ).map(({ id }) => id) } diff --git a/src/main/native-chat/agent-session-wire/structured-agent-session-options-read.test.ts b/src/main/native-chat/agent-session-wire/structured-agent-session-options-read.test.ts new file mode 100644 index 00000000000..5657fcfb4d8 --- /dev/null +++ b/src/main/native-chat/agent-session-wire/structured-agent-session-options-read.test.ts @@ -0,0 +1,48 @@ +// A chat's options at rest come from the host catalog without waiting on a listing. + +import { describe, expect, it, vi } from 'vitest' +import type { AgentSessionRecord } from '../../../shared/agent-session-record' +import { createAgentModelCatalogService } from '../agent-model-catalog/agent-model-catalog-service' +import { + AgentModelCatalogStore, + type AgentModelCatalogSuccess +} from '../agent-model-catalog/agent-model-catalog-store' +import type { StructuredAgentSessionMutationContext } from './structured-agent-session-host-mutations' +import { readStructuredAgentSessionOptions } from './structured-agent-session-options-read' + +const SESSION = 'session-1' + +function restingRecord(): AgentSessionRecord { + // oxlint-disable-next-line typescript/consistent-type-assertions -- SAFETY: the resting read and the catalog key touch only these fields. + return { + provider: 'codex', + accountHome: { variable: 'CODEX_HOME', path: '/homes/a' }, + location: { wslDistro: null }, + options: {} + } as unknown as AgentSessionRecord +} + +describe('options at rest', () => { + it('answers while the first catalog listing is still running', async () => { + const record = restingRecord() + const probe = vi.fn(() => new Promise(() => {})) + const modelCatalog = createAgentModelCatalogService({ + store: new AgentModelCatalogStore(), + getRecord: () => record, + resolveAccountHome: async () => ({ variable: 'CODEX_HOME', path: '/homes/a' }), + probes: { codex: probe } + }) + const resting = { child: null, params: { provider: 'codex' } } + // oxlint-disable-next-line typescript/consistent-type-assertions -- SAFETY: the resting read touches only these members. + const context = { + deps: { adapter: {}, store: { getRecord: () => record }, modelCatalog }, + serialize: (_sessionId: string, task: () => Promise) => task(), + openConversation: async () => resting, + conversation: async () => resting + } as unknown as StructuredAgentSessionMutationContext + + const result = await readStructuredAgentSessionOptions(context, SESSION) + expect(probe).toHaveBeenCalledTimes(1) + expect(result.models).toEqual([]) + }) +}) diff --git a/src/main/native-chat/agent-session-wire/structured-agent-session-start-failure-row.ts b/src/main/native-chat/agent-session-wire/structured-agent-session-start-failure-row.ts index 12fa05710f0..9206f95fde0 100644 --- a/src/main/native-chat/agent-session-wire/structured-agent-session-start-failure-row.ts +++ b/src/main/native-chat/agent-session-wire/structured-agent-session-start-failure-row.ts @@ -42,9 +42,9 @@ export function hasStructuredAgentSessionStartFailureRow( } /** - * A start the delivery loop needed and did not get: the start's row, and every queued message - * rejected with the same words. Writes nothing when nothing is still queued: a start whose - * messages Stop withdrew did not fail anyone. + * A start the delivery loop needed and did not get: every queued message rejected with the same + * words, then the start's row. Writes nothing when nothing is still queued: a start whose messages + * Stop withdrew did not fail anyone. */ export async function recordStructuredAgentSessionStartFailure( session: Pick & { fence: number }, @@ -59,11 +59,8 @@ export async function recordStructuredAgentSessionStartFailure( settlementId: `start-failure:${startKey}`, fence: session.fence, recovered: true, - mutations: [structuredAgentSessionStartFailureRow(startKey, failure)] - }) - await session.journal.rejectQueuedSubmissions(session.fence, { - reason: failure.reason, - rejection: failure.rejection + mutations: [structuredAgentSessionStartFailureRow(startKey, failure)], + rejectsQueued: { reason: failure.reason, rejection: failure.rejection } }) } diff --git a/src/main/providers/agent-foreground-process-pi.test.ts b/src/main/providers/agent-foreground-process-pi.test.ts index 691ee63e748..7ae8ec3f3fc 100644 --- a/src/main/providers/agent-foreground-process-pi.test.ts +++ b/src/main/providers/agent-foreground-process-pi.test.ts @@ -33,60 +33,64 @@ describe('Pi Windows foreground recognition', () => { } }) - it('recognizes the npm entrypoint within the active ConPTY', async () => { - const rows = [ - { - pid: 100, - ppid: 99, - name: 'bash.exe', - commandLine: '"C:\\Program Files\\Git\\usr\\bin\\bash.exe"' - }, - { - pid: 101, - ppid: 100, - name: 'node.exe', - commandLine: - 'node.exe C:\\Users\\dev\\AppData\\Roaming\\npm\\node_modules\\@earendil-works\\pi-coding-agent\\dist\\cli.js' - } - ] - getAllProcessesMock.mockImplementation((cb: (snapshot: unknown) => void) => { - cb(withSelf(rows)) - }) - const readWindowsConsoleAttachedProcessIds = vi.fn(async () => new Set([100, 101])) - - await expect( - resolveAgentForegroundProcessWithAvailability(100, 'node.exe', { - fresh: true, - readWindowsConsoleAttachedProcessIds + it.each(['cli.js', String.raw`bundle\cli.js`])( + 'recognizes the npm entrypoint dist/%s within the active ConPTY', + async (entrypoint) => { + const rows = [ + { + pid: 100, + ppid: 99, + name: 'bash.exe', + commandLine: '"C:\\Program Files\\Git\\usr\\bin\\bash.exe"' + }, + { + pid: 101, + ppid: 100, + name: 'node.exe', + commandLine: `node.exe C:\\Users\\dev\\AppData\\Roaming\\npm\\node_modules\\@earendil-works\\pi-coding-agent\\dist\\${entrypoint}` + } + ] + getAllProcessesMock.mockImplementation((cb: (snapshot: unknown) => void) => { + cb(withSelf(rows)) }) - ).resolves.toEqual({ available: true, processName: 'pi', processId: 101 }) - expect(readWindowsConsoleAttachedProcessIds).toHaveBeenCalledTimes(1) - }) + const readWindowsConsoleAttachedProcessIds = vi.fn(async () => new Set([100, 101])) - it('anchors a collapsed omp name to the omp pid, not the embedded pi leaf', async () => { - // Pi restarts under a live OMP; an anchor on pi's pid would read that as - // OMP's exit and fire a false "agent done" when the next snapshot degrades. - const rows = [ - { pid: 100, ppid: 99, name: 'powershell.exe', commandLine: 'powershell.exe' }, - { pid: 101, ppid: 100, name: 'omp.exe', commandLine: 'omp' }, - { - pid: 102, - ppid: 101, - name: 'node.exe', - commandLine: - 'node.exe C:\\npm\\node_modules\\@earendil-works\\pi-coding-agent\\dist\\cli.js' - } - ] - getAllProcessesMock.mockImplementation((cb: (snapshot: unknown) => void) => { - cb(withSelf(rows)) - }) - const readWindowsConsoleAttachedProcessIds = vi.fn(async () => new Set([100, 101, 102])) + await expect( + resolveAgentForegroundProcessWithAvailability(100, 'node.exe', { + fresh: true, + readWindowsConsoleAttachedProcessIds + }) + ).resolves.toEqual({ available: true, processName: 'pi', processId: 101 }) + expect(readWindowsConsoleAttachedProcessIds).toHaveBeenCalledTimes(1) + } + ) - await expect( - resolveAgentForegroundProcessWithAvailability(100, 'powershell.exe', { - fresh: true, - readWindowsConsoleAttachedProcessIds + it.each(['cli.js', String.raw`bundle\cli.js`])( + 'anchors a collapsed omp name to the omp pid, not the embedded pi leaf at dist/%s', + async (entrypoint) => { + // Pi restarts under a live OMP; an anchor on pi's pid would read that as + // OMP's exit and fire a false "agent done" when the next snapshot degrades. + const rows = [ + { pid: 100, ppid: 99, name: 'powershell.exe', commandLine: 'powershell.exe' }, + { pid: 101, ppid: 100, name: 'omp.exe', commandLine: 'omp' }, + { + pid: 102, + ppid: 101, + name: 'node.exe', + commandLine: `node.exe C:\\npm\\node_modules\\@earendil-works\\pi-coding-agent\\dist\\${entrypoint}` + } + ] + getAllProcessesMock.mockImplementation((cb: (snapshot: unknown) => void) => { + cb(withSelf(rows)) }) - ).resolves.toEqual({ available: true, processName: 'omp', processId: 101 }) - }) + const readWindowsConsoleAttachedProcessIds = vi.fn(async () => new Set([100, 101, 102])) + + await expect( + resolveAgentForegroundProcessWithAvailability(100, 'powershell.exe', { + fresh: true, + readWindowsConsoleAttachedProcessIds + }) + ).resolves.toEqual({ available: true, processName: 'omp', processId: 101 }) + } + ) }) diff --git a/src/main/rate-limits/cursor-fetcher.test.ts b/src/main/rate-limits/cursor-fetcher.test.ts index bb07fb90b05..a075c37db99 100644 --- a/src/main/rate-limits/cursor-fetcher.test.ts +++ b/src/main/rate-limits/cursor-fetcher.test.ts @@ -78,6 +78,7 @@ describe('fetchCursorRateLimits', () => { await fetchCursorRateLimits({ authReadResult: session() }) const [url, init] = netFetchMock.mock.calls[0] ?? [] expect(url).toBe('https://cursor.com/api/usage-summary') + expect(init?.credentials).toBe('omit') expect(init?.headers).toMatchObject({ Origin: 'https://cursor.com', Referer: 'https://cursor.com/dashboard', diff --git a/src/main/rate-limits/cursor-fetcher.ts b/src/main/rate-limits/cursor-fetcher.ts index 82bbf1c0e48..d6cbe09cfad 100644 --- a/src/main/rate-limits/cursor-fetcher.ts +++ b/src/main/rate-limits/cursor-fetcher.ts @@ -84,6 +84,8 @@ async function fetchDashboardJson( const res = await net.fetch(url, { // Keep dashboard redirects visible without forwarding session credentials. redirect: 'manual', + // The selected account's Cookie must not be replaced by Electron's session jar. + credentials: 'omit', headers: requestHeaders(session), signal: requestSignal }) diff --git a/src/main/runtime/rpc/methods/structured-agent-session-options-read.test.ts b/src/main/runtime/rpc/methods/structured-agent-session-options-read.test.ts index 77d44ece745..459ae5c75a4 100644 --- a/src/main/runtime/rpc/methods/structured-agent-session-options-read.test.ts +++ b/src/main/runtime/rpc/methods/structured-agent-session-options-read.test.ts @@ -56,4 +56,66 @@ describe('agentSession.modelCatalog', () => { await call('agentSession.modelCatalog', { agent: 'claude' }, STRUCTURED_CLIENT) expect(read).toHaveBeenCalledWith({ agent: 'claude' }) }) + + it('passes a wait for the listing through to the catalog', async () => { + await call( + 'agentSession.modelCatalog', + { agent: 'codex', sessionId: SESSION, waitForListing: true }, + STRUCTURED_CLIENT + ) + expect(read).toHaveBeenCalledWith({ agent: 'codex', sessionId: SESSION, waitForListing: true }) + }) +}) + +describe('agentSession.modelCatalog before anything built the host', () => { + const read = vi.fn(async () => ({ + origin: 'probe' as const, + models: [{ id: 'gpt-host', label: 'GPT Host', isDefault: true, efforts: [] }], + fetchedAt: 1 + })) + const installHost = vi.fn(async () => { + setStructuredAgentSessionHost(Object.assign(hostStub(), { deps: { modelCatalog: { read } } })) + }) + + beforeEach(() => { + read.mockClear() + installHost.mockClear() + clearStructuredHostStub() + }) + + // A new chat's picker reads before its create lands; on a host with no saved chats nothing else + // has built the host yet, and a refusal here left the picker on the client's built-in list. + it('builds the host for a structured chat and answers from its catalog', async () => { + const reply = await call( + 'agentSession.modelCatalog', + { agent: 'codex', sessionId: SESSION }, + STRUCTURED_CLIENT, + { ensureStructuredAgentSessionHost: installHost } + ) + expect(installHost).toHaveBeenCalledTimes(1) + expect(reply).toMatchObject({ ok: true, result: { origin: 'probe' } }) + expect(read).toHaveBeenCalledWith({ agent: 'codex', sessionId: SESSION }) + }) + + // Terminal-backed chat reads with no session: a host that runs no structured chat keeps its + // journal closed, and the read falls back to the CLI listing. + it('does not build the host for a read that names no session', async () => { + const reply = await call('agentSession.modelCatalog', { agent: 'codex' }, STRUCTURED_CLIENT, { + ensureStructuredAgentSessionHost: installHost + }) + expect(installHost).not.toHaveBeenCalled() + expect(reply).toMatchObject({ ok: false }) + expect(read).not.toHaveBeenCalled() + }) + + it('does not build the host for a client that cannot read structured sessions', async () => { + const reply = await call( + 'agentSession.modelCatalog', + { agent: 'codex', sessionId: SESSION }, + { clientKind: 'runtime', clientCapabilities: [] }, + { ensureStructuredAgentSessionHost: installHost } + ) + expect(installHost).not.toHaveBeenCalled() + expect(reply).toMatchObject({ ok: false }) + }) }) diff --git a/src/main/runtime/rpc/methods/structured-agent-session-options-read.ts b/src/main/runtime/rpc/methods/structured-agent-session-options-read.ts index 19e93689577..3cd86c876ad 100644 --- a/src/main/runtime/rpc/methods/structured-agent-session-options-read.ts +++ b/src/main/runtime/rpc/methods/structured-agent-session-options-read.ts @@ -11,7 +11,7 @@ import { defineMethod } from '../core' import { requireInstalledStructuredHost, - requireStructuredHost as requireHost + requireStructuredHost } from './structured-agent-session-gate' import { ModelCatalogParams, OptionsParams } from './structured-agent-session-schemas' @@ -25,8 +25,14 @@ export const STRUCTURED_AGENT_SESSION_OPTIONS_READ_METHODS = [ defineMethod({ name: 'agentSession.modelCatalog', params: ModelCatalogParams, + // A structured chat's read names its session and builds the host, since it may come first; + // terminal-backed chat's session-less read must not open the journal where none runs. handler: async ({ worktree, ...params }, ctx) => { - const catalog = requireHost(ctx).deps.modelCatalog + const host = + params.sessionId === undefined + ? requireStructuredHost(ctx) + : await requireInstalledStructuredHost(ctx) + const catalog = host.deps.modelCatalog if (!catalog) { return { origin: 'unknown' as const } } diff --git a/src/main/worker-thread-request-queue.test.ts b/src/main/worker-thread-request-queue.test.ts index 4bf2188dc54..2046267396c 100644 --- a/src/main/worker-thread-request-queue.test.ts +++ b/src/main/worker-thread-request-queue.test.ts @@ -92,9 +92,10 @@ function makeQueue( function send( queue: WorkerThreadRequestQueue, - label: string + label: string, + owner?: { readonly signal: AbortSignal } ): Promise { - return queue.dispatch((id) => ({ id, label }), TIMEOUT_MS) + return queue.dispatch((id) => ({ id, label }), TIMEOUT_MS, undefined, owner) } /** Resolve to the response or to the rejection, so a test can assert on either. */ @@ -328,6 +329,161 @@ describe('WorkerThreadRequestQueue', () => { expect(await behind).toMatchObject({ label: 'd' }) }) + it('keeps one owner fault budget across idle batches and refuses a fourth worker', async () => { + const workers: FakeWorker[] = [] + const queue = makeQueue(workers) + const owner = { signal: new AbortController().signal } + try { + for (let fault = 0; fault < 3; fault++) { + const call = settle(send(queue, `owned-${fault}`, owner)) + workers.at(-1)?.emit('error', new Error(`fault-${fault}`)) + expect(await call).toMatchObject({ message: `fault-${fault}` }) + } + const refused = settle(send(queue, 'fourth', owner)) + expect(workers).toHaveLength(3) + await expect(refused).resolves.toMatchObject({ message: 'crashed repeatedly (fault-2)' }) + const nextScan = send(queue, 'new-scan', { signal: new AbortController().signal }) + expect(workers).toHaveLength(4) + workers[3].respond() + await expect(nextScan).resolves.toMatchObject({ label: 'new-scan' }) + } finally { + queue.dispose() + } + }) + + it('resets only the successful owner and does not count idle worker exits', async () => { + const workers: FakeWorker[] = [] + const queue = makeQueue(workers) + const a = { signal: new AbortController().signal } + const b = { signal: new AbortController().signal } + try { + const initial = send(queue, 'initial', a) + workers[0].respond() + await initial + workers[0].emit('exit', 17) + for (let fault = 0; fault < 2; fault++) { + const call = settle(send(queue, `a-${fault}`, a)) + workers.at(-1)?.emit('error', new Error(`a-fault-${fault}`)) + await call + } + const peer = send(queue, 'b-success', b) + workers.at(-1)?.respond() + await peer + const third = settle(send(queue, 'a-third', a)) + workers.at(-1)?.emit('error', new Error('a-third-fault')) + await third + const refused = settle(send(queue, 'a-fourth', a)) + expect(workers).toHaveLength(4) + await expect(refused).resolves.toMatchObject({ + message: 'crashed repeatedly (a-third-fault)' + }) + const healthy = send(queue, 'b-still-healthy', b) + workers.at(-1)?.respond() + await expect(healthy).resolves.toMatchObject({ label: 'b-still-healthy' }) + } finally { + queue.dispose() + } + }) + + it('clears a successful owner budget while preserving another owner failure count', async () => { + const workers: FakeWorker[] = [] + const queue = makeQueue(workers) + const a = { signal: new AbortController().signal } + const b = { signal: new AbortController().signal } + try { + for (const owner of [a, b]) { + for (let fault = 0; fault < 2; fault++) { + const call = settle(send(queue, 'failure', owner)) + workers.at(-1)?.emit('error', new Error('failure')) + await call + } + } + const successful = send(queue, 'a-success', a) + workers.at(-1)?.respond() + await successful + const thirdB = settle(send(queue, 'b-third', b)) + workers.at(-1)?.emit('error', new Error('b-third-fault')) + await thirdB + for (let fault = 0; fault < 2; fault++) { + const call = settle(send(queue, 'a-new-failure', a)) + workers.at(-1)?.emit('error', new Error('a-new-failure')) + await call + } + const stillAllowed = send(queue, 'a-allowed', a) + expect(workers).toHaveLength(8) + workers.at(-1)?.respond() + await expect(stillAllowed).resolves.toMatchObject({ label: 'a-allowed' }) + await expect(settle(send(queue, 'b-refused', b))).resolves.toMatchObject({ + message: 'crashed repeatedly (b-third-fault)' + }) + expect(workers).toHaveLength(8) + } finally { + queue.dispose() + } + }) + + it('drains only the failing owner while queued peers keep their FIFO order', async () => { + const workers: FakeWorker[] = [] + const queue = makeQueue(workers) + const a = { signal: new AbortController().signal } + const b = { signal: new AbortController().signal } + try { + const failed = ['a1', 'a2', 'a3', 'a4'].map((label) => settle(send(queue, label, a))) + const peer = settle(send(queue, 'peer', b)) + const ordinary = settle(send(queue, 'ordinary')) + for (let fault = 0; fault < 3; fault++) { + workers.at(-1)?.emit('error', new Error(`fault-${fault}`)) + } + expect(workers).toHaveLength(4) + expect(labels(workers[3])).toEqual(['peer']) + expect((await Promise.all(failed)).at(-1)).toMatchObject({ + message: 'crashed repeatedly (fault-2)' + }) + workers[3].respond() + await expect(peer).resolves.toMatchObject({ label: 'peer' }) + expect(labels(workers[3])).toEqual(['peer', 'ordinary']) + workers[3].respond() + await expect(ordinary).resolves.toMatchObject({ label: 'ordinary' }) + } finally { + queue.dispose() + } + }) + + it('keeps retirement refusals out of the owner failure budget', async () => { + vi.useFakeTimers() + const workers: FakeWorker[] = [] + let finish: (code: number) => void = () => {} + const first = new FakeWorker() + first.exit = new Promise((resolve) => { + finish = resolve + }) + const queue = makeQueue(workers, { + awaitRetirement: true, + makeWorker: () => (workers.length === 0 ? first : new FakeWorker()) + }) + const owner = { signal: new AbortController().signal } + try { + const failed = settle(send(queue, 'first', owner)) + first.emit('error', new Error('first-fault')) + await failed + for (let attempt = 0; attempt < 4; attempt++) { + await expect(settle(send(queue, 'not-executed', owner))).resolves.toMatchObject({ + message: 'unavailable: previous worker still exiting' + }) + } + expect(workers).toHaveLength(1) + finish(1) + await vi.advanceTimersByTimeAsync(0) + const recovered = send(queue, 'recovered', owner) + expect(workers).toHaveLength(2) + workers[1].respond() + await expect(recovered).resolves.toMatchObject({ label: 'recovered' }) + } finally { + finish(1) + queue.dispose() + } + }) + describe('awaitRetirement', () => { function stalledExit(): { worker: FakeWorker; finish: (code: number) => void } { const worker = new FakeWorker() diff --git a/src/main/worker-thread-request-queue.ts b/src/main/worker-thread-request-queue.ts index 39561fa80a4..6793d597e27 100644 --- a/src/main/worker-thread-request-queue.ts +++ b/src/main/worker-thread-request-queue.ts @@ -42,6 +42,10 @@ export type WorkerThreadRequestQueueOptions = { awaitRetirement?: boolean } +export type WorkerThreadRequestOwner = { readonly signal: AbortSignal } + +type OwnerFailures = { consecutiveDeaths: number; refused: Error | null } + type PendingCall = { request: TRequest timeoutMs: number @@ -49,6 +53,7 @@ type PendingCall = { reject: (error: Error) => void timer: NodeJS.Timeout | null signal?: AbortSignal + owner?: WorkerThreadRequestOwner cleanupAbort: () => void } @@ -59,6 +64,7 @@ export class WorkerThreadRequestQueue< private active: PendingCall | null = null private queue: PendingCall[] = [] private consecutiveDeaths = 0 + private readonly ownerFailures = new WeakMap() private nextId = 1 private disposed = false private readonly host: LazyWorkerThreadHost @@ -85,7 +91,8 @@ export class WorkerThreadRequestQueue< dispatch( buildRequest: (id: number) => TRequest, timeoutMs: number, - signal?: AbortSignal + signal?: AbortSignal, + owner?: WorkerThreadRequestOwner ): Promise { return new Promise((resolve, reject) => { if (this.disposed) { @@ -96,6 +103,11 @@ export class WorkerThreadRequestQueue< reject(signal.reason ?? new Error('Worker request aborted')) return } + const refused = owner && this.ownerFailures.get(owner)?.refused + if (refused) { + reject(refused) + return + } // Built before the cap check so a rejection can name the dropped work; // the id it burns is only a correlation token, so a gap costs nothing. const request = buildRequest(this.nextId++) @@ -106,7 +118,7 @@ export class WorkerThreadRequestQueue< } // A fresh burst from full idle starts new work: clear any death count // carried from a prior burst so the respawn cap can't drain it early. - if (!this.active && this.queue.length === 0) { + if (!owner && !this.active && this.queue.length === 0) { this.consecutiveDeaths = 0 } const call: PendingCall = { @@ -116,6 +128,7 @@ export class WorkerThreadRequestQueue< reject, timer: null, signal, + owner, cleanupAbort: () => signal?.removeEventListener('abort', abort) } const abort = (): void => { @@ -201,7 +214,11 @@ export class WorkerThreadRequestQueue< this.armDeadline(call) return } - this.consecutiveDeaths = 0 + if (call.owner) { + this.failuresFor(call.owner).consecutiveDeaths = 0 + } else { + this.consecutiveDeaths = 0 + } this.settle(call, () => call.resolve(response)) this.afterSettle() } @@ -226,12 +243,16 @@ export class WorkerThreadRequestQueue< private onWorkerFault(error: Error): void { const failed = this.active this.host.destroy() - this.consecutiveDeaths++ - if (failed) { - this.settle(failed, () => failed.reject(error)) + if (!failed) { + this.pump() + return } - if (this.consecutiveDeaths >= this.options.maxConsecutiveDeaths) { - this.drainQueueAfterCrashLoop(error) + const deaths = failed.owner + ? ++this.failuresFor(failed.owner).consecutiveDeaths + : ++this.consecutiveDeaths + this.settle(failed, () => failed.reject(error)) + if (deaths >= this.options.maxConsecutiveDeaths) { + this.drainQueueAfterCrashLoop(error, failed.owner) return } if (this.queue.length > 0) { @@ -239,14 +260,28 @@ export class WorkerThreadRequestQueue< } } - private drainQueueAfterCrashLoop(error: Error): void { - const pending = this.queue - this.queue = [] - this.consecutiveDeaths = 0 + private drainQueueAfterCrashLoop(error: Error, owner?: WorkerThreadRequestOwner): void { + const pending = this.queue.filter((call) => call.owner === owner) + this.queue = this.queue.filter((call) => call.owner !== owner) const drainError = new Error(this.options.describeCrashLoop(error.message)) + if (owner) { + this.failuresFor(owner).refused = drainError + } else { + this.consecutiveDeaths = 0 + } for (const call of pending) { this.settle(call, () => call.reject(drainError)) } + this.afterSettle() + } + + private failuresFor(owner: WorkerThreadRequestOwner): OwnerFailures { + let failures = this.ownerFailures.get(owner) + if (!failures) { + failures = { consecutiveDeaths: 0, refused: null } + this.ownerFailures.set(owner, failures) + } + return failures } private failQueuedAsUnavailable(): void { diff --git a/src/renderer/src/components/mobile/mobile-platform-copy.ts b/src/renderer/src/components/mobile/mobile-platform-copy.ts index f33f9a42196..45149f90858 100644 --- a/src/renderer/src/components/mobile/mobile-platform-copy.ts +++ b/src/renderer/src/components/mobile/mobile-platform-copy.ts @@ -22,7 +22,7 @@ const IOS_CHANNEL_COPY: Record = { const ANDROID_COPY: InstallCopy = { ctaLabel: 'Download APK', - url: 'https://github.com/stablyai/orca/releases/download/mobile-android-v0.0.48/app-release.apk' + url: 'https://github.com/stablyai/orca/releases/download/mobile-android-v0.0.52/app-release.apk' } export function getInstallCopy(platform: Platform, iosChannel: IosChannel): InstallCopy { diff --git a/src/renderer/src/components/native-chat/NativeChatMessageList.provider-retry-runs.test.tsx b/src/renderer/src/components/native-chat/NativeChatMessageList.provider-retry-runs.test.tsx index eb71e3e6b54..db59afdb19d 100644 --- a/src/renderer/src/components/native-chat/NativeChatMessageList.provider-retry-runs.test.tsx +++ b/src/renderer/src/components/native-chat/NativeChatMessageList.provider-retry-runs.test.tsx @@ -10,6 +10,8 @@ import { NativeChatMessageList } from './NativeChatMessageList' import { installNativeChatMessageListTestViewport } from './native-chat-message-list-test-viewport' import { projectStructuredAgentSessionMessages } from './structured-agent-session-message-projection' +const NO_CARDS: readonly string[] = [] + let restoreViewport = (): void => {} beforeAll(() => { restoreViewport = installNativeChatMessageListTestViewport() @@ -42,7 +44,7 @@ function transcript(items: AgentJournalRenderItem[]) { return ( {} +beforeAll(() => { + restoreViewport = installNativeChatMessageListTestViewport() +}) +afterAll(() => restoreViewport()) +afterEach(cleanup) + +type Phase = 'in flight' | 'running' | 'done' + +const SEED = agentJournalSubmissionKey('seed') +const NEW = agentJournalSubmissionKey('new') + +function said(role: 'user' | 'assistant', text: string): AgentJournalItemBody { + return { kind: 'message', role, blocks: [{ type: 'text', text }] } +} + +function journal(phase: Phase, scoped: boolean): AgentJournalRenderItem[] { + const thread: AgentJournalTurnScope = { kind: 'thread' } + const row = ( + itemId: string, + body: AgentJournalItemBody, + turnScope: AgentJournalTurnScope + ): AgentJournalRenderItem => ({ + itemId, + revision: 0, + sequence: rows.length + 1, + observedAt: 1000 + rows.length, + body, + ...(scoped ? { turnScope } : {}) + }) + const rows: AgentJournalRenderItem[] = [] + rows.push(row(SEED, said('user', 'SEED PROMPT'), thread)) + rows.push( + row( + 'turn-seed', + { kind: 'turn', turnId: 'turn-seed', state: 'completed', userItemId: SEED }, + thread + ) + ) + rows.push( + row('seed-answer', said('assistant', 'SEED OK'), { kind: 'turn', turnItemId: 'turn-seed' }) + ) + rows.push(row(NEW, said('user', 'NEW PROMPT'), thread)) + if (phase !== 'in flight') { + rows.push( + row( + 'turn-new', + { + kind: 'turn', + turnId: 'turn-new', + state: phase === 'done' ? 'completed' : 'running', + userItemId: NEW + }, + thread + ) + ) + } + if (phase === 'done') { + rows.push( + row('new-answer', said('assistant', 'NEW OK'), { kind: 'turn', turnItemId: 'turn-new' }) + ) + } + return rows +} + +function submission( + id: string, + dispatchState: AgentJournalSubmission['dispatchState'] +): AgentJournalSubmission { + return { + clientMessageId: id, + fence: 1, + payloadFingerprint: id, + dispatchState, + providerItemId: null, + reason: null, + submittedAt: 1, + resolvedAt: null + } +} + +/** The drawn sequence of the prompts, bars and live activity line, top to bottom. */ +function drawn(container: HTMLElement): string[] { + const out: string[] = [] + for (const element of container.querySelectorAll('*')) { + if (element.hasAttribute('data-native-chat-turn-activity')) { + out.push('ACTIVITY') + } else if (element.hasAttribute('data-native-chat-turn-status')) { + out.push( + element.getAttribute('data-native-chat-turn-status') === 'active' ? 'WORKING' : 'WORKED' + ) + } else if (element.childElementCount === 0 && (element.textContent ?? '').endsWith('PROMPT')) { + out.push(element.textContent ?? '') + } + } + return out +} + +describe('a message the host rejected after a crash, with no outbox entry left', () => { + const LOST = agentJournalSubmissionKey('lost') + + function hostList(phase: Phase) { + const items = journal(phase, true) + const newIndex = items.findIndex((item) => item.itemId === NEW) + // Recorded between the seed's answer and the newer prompt; the next open rejected it. + items.splice(newIndex, 0, { + itemId: LOST, + revision: 0, + sequence: items[newIndex - 1]!.sequence, + sequenceIndex: 1, + observedAt: 1002.5, + body: said('user', 'LOST PROMPT'), + turnScope: { kind: 'thread' } + }) + const submissions: AgentJournalSubmission[] = [ + submission('seed', 'accepted'), + { + ...submission('lost', 'rejected'), + reason: DISPATCH_REJECTED_HOST_RESTARTED, + rejection: { kind: 'hostRestarted' } + }, + submission('new', phase === 'done' ? 'accepted' : 'pending') + ] + const settledTurns: NativeChatSettledTurns = new Map([ + [SEED, { startedAt: 1, workedSeconds: 3 }], + ...(phase === 'done' ? [[NEW, { startedAt: 2, workedSeconds: 5 }] as const] : []) + ]) + return ( + + ) + } + + it('draws it where it was recorded, muted with its reason and no control, outside the newer turn', () => { + const { container, rerender } = render(hostList('running')) + expect(drawn(container)).toEqual([ + 'SEED PROMPT', + 'WORKED', + 'LOST PROMPT', + 'NEW PROMPT', + 'WORKING', + 'ACTIVITY' + ]) + expect(container.textContent).toContain('Orca restarted before this message was sent.') + expect(container.querySelector('button[aria-label="Retry"]')).toBeNull() + expect( + [...container.querySelectorAll('button')].map((button) => button.textContent) + ).not.toContain('Retry') + rerender(hostList('done')) + expect(drawn(container)).toEqual([ + 'SEED PROMPT', + 'WORKED', + 'LOST PROMPT', + 'NEW PROMPT', + 'WORKED' + ]) + expect(container.textContent).toContain('Worked for 5s') + }) +}) diff --git a/src/renderer/src/components/native-chat/NativeChatMessageList.stream-render.perf.test.tsx b/src/renderer/src/components/native-chat/NativeChatMessageList.stream-render.perf.test.tsx index 74df2b75bd6..7efb254c619 100644 --- a/src/renderer/src/components/native-chat/NativeChatMessageList.stream-render.perf.test.tsx +++ b/src/renderer/src/components/native-chat/NativeChatMessageList.stream-render.perf.test.tsx @@ -143,7 +143,9 @@ describe('native chat transcript re-render cost during a streaming turn', () => } const view = (items: AgentJournalRenderItem[]) => ( Promise.resolve('exhausted' as const) function Transcript({ items }: { items: AgentJournalRenderItem[] }) { - const messages = useStructuredAgentSessionMessages(items, EMPTY, EMPTY) + const messages = useStructuredAgentSessionMessages(items, EMPTY, EMPTY, EMPTY) const session: NativeChatLiveSession = { messages, status: 'working', diff --git a/src/renderer/src/components/native-chat/NativeChatMessageList.turn-membership.test.tsx b/src/renderer/src/components/native-chat/NativeChatMessageList.turn-membership.test.tsx index b01675c63a7..29f977046f6 100644 --- a/src/renderer/src/components/native-chat/NativeChatMessageList.turn-membership.test.tsx +++ b/src/renderer/src/components/native-chat/NativeChatMessageList.turn-membership.test.tsx @@ -88,7 +88,9 @@ function journalList( return ( ({ toast: { error: vi.fn() } })) + +vi.mock('@/i18n/i18n', () => ({ + translate: (_key: string, fallback: string, values?: Record) => + values + ? Object.entries(values).reduce( + (text, [name, value]) => text.replaceAll(`{{${name}}}`, String(value)), + fallback + ) + : fallback +})) + +import { TooltipProvider } from '@/components/ui/tooltip' +import { NativeChatSessionOptionPickers } from './NativeChatSessionOptionPickers' +import type { NativeChatOptionPickerRequest } from './native-chat-composer-types' + +function model(pending: boolean): SessionOptionDescriptor { + return { + id: 'model', + label: 'Model', + category: 'model', + kind: { + type: 'select', + currentValue: 'opus', + choices: [ + { value: 'opus', label: 'Opus 4.8' }, + { value: 'sonnet', label: 'Sonnet 5' } + ] + }, + valueSource: 'applied', + transport: 'agent-session', + settable: !pending, + ...(pending ? { choicesPending: true as const } : {}) + } +} + +const surface: SessionOptionsSurface = { + getSnapshot: () => [], + setOption: vi.fn(async () => ({ snapshot: [] })), + invokeAction: vi.fn(async () => ({ snapshot: [] })), + subscribe: () => () => {} +} + +function view(pending: boolean, request: NativeChatOptionPickerRequest | null): React.JSX.Element { + return ( + +