From 0ff258490c262cb1059954828b8ec396d54426df Mon Sep 17 00:00:00 2001 From: Merge Sim Date: Tue, 8 Sep 2026 00:19:17 -0700 Subject: [PATCH] Implement request-owned Linear page acquisition and recovery --- src/cli/handlers/linear-list-issues.ts | 19 ++ src/cli/linear-format.ts | 6 +- src/cli/specs/linear-mcp.ts | 1 + src/main/linear/client.ts | 1 + src/main/linear/issue-context-client.ts | 8 +- .../linear/linear-account-read-lifetime.ts | 29 ++ src/main/linear/linear-request-concurrency.ts | 56 +++- src/main/linear/linear-token-store.ts | 3 + src/main/linear/linear-workspace-registry.ts | 2 + src/main/linear/mcp-issue-list-acquisition.ts | 148 +++++++++ src/main/linear/mcp-issue-list-admission.ts | 82 +++++ .../linear/mcp-issue-list-lifetime.test.ts | 94 ++++++ src/main/linear/mcp-issue-list-lifetime.ts | 90 ++++++ .../linear/mcp-issue-list-page-owned.test.ts | 264 +++++++++++++++++ src/main/linear/mcp-issue-list-pages.ts | 280 ++++++++++++++++++ src/main/linear/mcp-issue-list-recovery.ts | 166 +++++++++++ src/main/linear/mcp-issue-list.test.ts | 45 ++- src/main/linear/mcp-issue-list.ts | 268 ++++------------- src/main/runtime/rpc/core.ts | 1 + .../runtime/rpc/dispatcher-stream-options.ts | 1 + src/main/runtime/rpc/dispatcher.ts | 24 +- .../runtime/rpc/linear-list-reply-budget.ts | 116 ++++++++ .../rpc/methods/linear-issue-list-method.ts | 14 +- src/main/runtime/rpc/methods/linear.test.ts | 28 +- .../runtime/rpc/rpc-streaming-dispatcher.ts | 28 +- src/main/runtime/rpc/transport.ts | 1 + src/main/runtime/rpc/unix-socket-transport.ts | 22 +- .../runtime/runtime-linear-read-commands.ts | 7 +- .../runtime-rpc/runtime-rpc-lifecycle.ts | 11 + .../runtime-rpc-request-admission.ts | 4 +- .../runtime-rpc-websocket-dispatch.ts | 22 +- src/main/ssh/linear-list-ssh-delivery.ts | 96 ++++++ src/main/ssh/ssh-channel-multiplexer.ts | 26 +- src/main/ssh/ssh-relay-session.ts | 8 +- .../ssh/ssh-remote-cli-host-passthrough.ts | 2 + .../ssh-remote-cli-interactive-commands.ts | 13 + src/main/ssh/ssh-remote-linear-list-issues.ts | 24 ++ src/main/ssh/ssh-remote-linear-output.ts | 3 + src/main/ssh/ssh-remote-linear-read-flags.ts | 1 + src/main/ssh/ssh-remote-orca-cli.ts | 31 +- src/shared/fetch-response-body.ts | 16 +- src/shared/linear/agent-access.ts | 11 +- src/shared/linear/mcp-issue-list.ts | 10 + src/shared/linear/workspace-types.ts | 1 + 44 files changed, 1786 insertions(+), 297 deletions(-) create mode 100644 src/main/linear/linear-account-read-lifetime.ts create mode 100644 src/main/linear/mcp-issue-list-acquisition.ts create mode 100644 src/main/linear/mcp-issue-list-admission.ts create mode 100644 src/main/linear/mcp-issue-list-lifetime.test.ts create mode 100644 src/main/linear/mcp-issue-list-lifetime.ts create mode 100644 src/main/linear/mcp-issue-list-page-owned.test.ts create mode 100644 src/main/linear/mcp-issue-list-pages.ts create mode 100644 src/main/linear/mcp-issue-list-recovery.ts create mode 100644 src/main/runtime/rpc/linear-list-reply-budget.ts create mode 100644 src/main/ssh/linear-list-ssh-delivery.ts create mode 100644 src/main/ssh/ssh-remote-cli-interactive-commands.ts diff --git a/src/cli/handlers/linear-list-issues.ts b/src/cli/handlers/linear-list-issues.ts index 51efa1b711e..e1ada10a55e 100644 --- a/src/cli/handlers/linear-list-issues.ts +++ b/src/cli/handlers/linear-list-issues.ts @@ -1,3 +1,4 @@ +import type { LinearConnectionStatus } from '../../shared/linear/workspace-types' import type { LinearMcpIssueListRequest, LinearMcpIssueListResult @@ -37,6 +38,24 @@ export const runLinearListIssues: CommandHandler = async ({ flags, client, json includeArchived: flags.get('include-archived') === true, workspaceId: getOptionalStringFlag(flags, 'workspace') } + const continuation = getOptionalStringFlag(flags, 'page-recovery') + if (continuation && (request.workspaceId !== 'all' || request.cursor)) { + throw new RuntimeClientError( + 'invalid_argument', + '--page-recovery requires --workspace all and cannot use --cursor' + ) + } + if (request.workspaceId === 'all') { + const status = await client.call('linear.status', {}) + if (status.result.mcpListPageRecoveryVersion === 1) { + request.pageRecovery = { version: 1, ...(continuation ? { continuation } : {}) } + } else if (continuation) { + throw new RuntimeClientError( + 'linear_list_concrete_workspace_required', + 'This runtime does not support page recovery; restart concrete workspaces and reconcile by issue ID.' + ) + } + } const response = await client.call('linear.mcpListIssues', request) if (!json) { printLinearMcpIssueListWarnings(response.result) diff --git a/src/cli/linear-format.ts b/src/cli/linear-format.ts index 57464f37e09..996ae237e1d 100644 --- a/src/cli/linear-format.ts +++ b/src/cli/linear-format.ts @@ -139,7 +139,11 @@ export function formatLinearMcpIssueList(result: LinearMcpIssueListResult): stri } export function printLinearMcpIssueListWarnings(result: LinearMcpIssueListResult): void { - if (result.meta.hasMore) { + if (result.meta.pageRecovery && result.meta.hasMore) { + console.error( + `warning: admitted batch; continue with --workspace all --page-recovery ${result.meta.pageRecovery.continuation}` + ) + } else if (result.meta.hasMore) { const workspaceHint = result.meta.nextCursor && result.meta.workspaceId !== 'all' && result.meta.workspaceId ? `; continue with --workspace ${result.meta.workspaceId}` diff --git a/src/cli/specs/linear-mcp.ts b/src/cli/specs/linear-mcp.ts index af7f654e0a3..e953a4a5701 100644 --- a/src/cli/specs/linear-mcp.ts +++ b/src/cli/specs/linear-mcp.ts @@ -53,6 +53,7 @@ export const LINEAR_MCP_COMMAND_SPECS: CommandSpec[] = [ 'query', 'state', 'cursor', + 'page-recovery', 'order-by', 'project', 'release', diff --git a/src/main/linear/client.ts b/src/main/linear/client.ts index 7e02938b54c..0f643294cc9 100644 --- a/src/main/linear/client.ts +++ b/src/main/linear/client.ts @@ -179,6 +179,7 @@ export function getStatus(): LinearConnectionStatus { return { connected: state.workspaces.length > 0, + mcpListPageRecoveryVersion: 1, viewer: activeWorkspace, workspaces: state.workspaces, activeWorkspaceId: state.activeWorkspaceId, diff --git a/src/main/linear/issue-context-client.ts b/src/main/linear/issue-context-client.ts index 7153133ca31..4a97ab797dc 100644 --- a/src/main/linear/issue-context-client.ts +++ b/src/main/linear/issue-context-client.ts @@ -181,12 +181,18 @@ export async function withLinearRead( try { return await read() } catch (error) { - if (isAuthError(error)) { + if ( + isAuthError(error) || + (error instanceof LinearAgentAccessError && error.code === 'linear_auth_expired') + ) { clearToken(entry.workspace.id) throw linearError('linear_auth_expired', 'Linear authentication expired.', { nextSteps: ['Reconnect Linear from Orca settings.'] }) } + if (error instanceof LinearAgentAccessError) { + throw error + } throw linearError(classifyLinearError(error), linearMessage(error)) } finally { release() diff --git a/src/main/linear/linear-account-read-lifetime.ts b/src/main/linear/linear-account-read-lifetime.ts new file mode 100644 index 00000000000..71b3b72b41e --- /dev/null +++ b/src/main/linear/linear-account-read-lifetime.ts @@ -0,0 +1,29 @@ +const reads = new Map>() + +export function registerLinearAccountRead(workspaceId: string): { + signal: AbortSignal + dispose: () => void +} { + const controller = new AbortController() + let active = reads.get(workspaceId) + if (!active) { + active = new Set() + reads.set(workspaceId, active) + } + active.add(controller) + return { + signal: controller.signal, + dispose: () => { + active.delete(controller) + if (active.size === 0) { + reads.delete(workspaceId) + } + } + } +} + +export function invalidateLinearAccountReads(workspaceId: string): void { + for (const controller of reads.get(workspaceId) ?? []) { + controller.abort(new Error('Linear account changed during the read.')) + } +} diff --git a/src/main/linear/linear-request-concurrency.ts b/src/main/linear/linear-request-concurrency.ts index b61d2c84aa7..95d12c1af06 100644 --- a/src/main/linear/linear-request-concurrency.ts +++ b/src/main/linear/linear-request-concurrency.ts @@ -1,25 +1,57 @@ -// ── Concurrency limiter — max 4 parallel Linear API calls ──────────── const MAX_CONCURRENT = 4 +const LIST_RESERVATION_BYTES = 1_179_648 +const LIST_RESERVATION_ALLOWANCE = 32 * 1024 * 1024 let running = 0 -const queue: (() => void)[] = [] +let reserved = 0 +const queue: { start: () => void; cancel: () => void }[] = [] -export function acquire(): Promise { +export function acquire(signal?: AbortSignal): Promise { + if (signal?.aborted) { + return Promise.reject(signal.reason) + } if (running < MAX_CONCURRENT) { running++ return Promise.resolve() } - return new Promise((resolve) => - queue.push(() => { - running++ - resolve() - }) - ) + return new Promise((resolve, reject) => { + const detach = (): void => signal?.removeEventListener('abort', waiter.cancel) + const waiter = { + start: (): void => { + detach() + running++ + resolve() + }, + cancel: (): void => { + const index = queue.indexOf(waiter) + if (index === -1) { + return + } + queue.splice(index, 1) + detach() + reject(signal?.reason) + } + } + queue.push(waiter) + signal?.addEventListener('abort', waiter.cancel, { once: true }) + }) } export function release(): void { running-- - const next = queue.shift() - if (next) { - next() + queue.shift()?.start() +} + +export function reserveLinearListing(): (() => void) | null { + if (reserved + LIST_RESERVATION_BYTES > LIST_RESERVATION_ALLOWANCE) { + return null + } + reserved += LIST_RESERVATION_BYTES + let released = false + return () => { + if (released) { + return + } + released = true + reserved -= LIST_RESERVATION_BYTES } } diff --git a/src/main/linear/linear-token-store.ts b/src/main/linear/linear-token-store.ts index 7a2bf892213..a79c34a3268 100644 --- a/src/main/linear/linear-token-store.ts +++ b/src/main/linear/linear-token-store.ts @@ -1,3 +1,4 @@ +import { invalidateLinearAccountReads } from './linear-account-read-lifetime' import { getSecretStore } from '../../shared/secret-store' import { existsSync, readFileSync, unlinkSync, writeFileSync } from 'node:fs' import { @@ -44,6 +45,7 @@ function writeEncryptedToken(path: string, apiKey: string): void { } export function saveWorkspaceToken(workspaceId: string, apiKey: string): void { + invalidateLinearAccountReads(workspaceId) ensureOrcaDir() if (workspaceId !== LEGACY_WORKSPACE_ID) { ensureWorkspaceTokenDir() @@ -93,6 +95,7 @@ export function loadToken(options: { force?: boolean; workspaceId?: string } = { } export function clearTokenFile(workspaceId: string): void { + invalidateLinearAccountReads(workspaceId) forgetCachedToken(workspaceId) try { unlinkSync(getWorkspaceTokenPath(workspaceId)) diff --git a/src/main/linear/linear-workspace-registry.ts b/src/main/linear/linear-workspace-registry.ts index 0a658a153ae..ffe1c1299bd 100644 --- a/src/main/linear/linear-workspace-registry.ts +++ b/src/main/linear/linear-workspace-registry.ts @@ -1,3 +1,4 @@ +import { invalidateLinearAccountReads } from './linear-account-read-lifetime' // ── Token + workspace storage ──────────────────────────────────────── // Why: tokens remain encrypted via safeStorage, while workspace metadata stays // plaintext so status checks can render connected accounts without decrypting @@ -197,6 +198,7 @@ export function upsertWorkspace( workspace: LinearWorkspace, options: { select?: boolean } = {} ): void { + invalidateLinearAccountReads(workspace.id) const file = getWorkspaceFile() const current = file.workspaces.find((entry) => entry.id === workspace.id) const credentialRevision = (current?.credentialRevision ?? 0) + 1 diff --git a/src/main/linear/mcp-issue-list-acquisition.ts b/src/main/linear/mcp-issue-list-acquisition.ts new file mode 100644 index 00000000000..3f5dc9d16f4 --- /dev/null +++ b/src/main/linear/mcp-issue-list-acquisition.ts @@ -0,0 +1,148 @@ +import type { LinearClientOptions } from '@linear/sdk' +import { z } from 'zod' +import { ISSUE_FIELDS } from './issue-context-raw' +import { linearError } from './issue-context-errors' +import { readFetchResponseBytesWithinLimit } from '../../shared/fetch-response-body' +import { assertJsonTextStructureWithinLimits } from '../../shared/json-text-structure-limit' + +const nullableText = z.string().nullish() +const named = z.object({ + id: nullableText, + name: nullableText, + color: nullableText, + type: nullableText +}) +const issue = z.object({ + id: z.string().min(1), + identifier: z.string().min(1), + title: z.string(), + url: z.string(), + description: nullableText, + priority: z.number().finite().nullish(), + estimate: z.number().finite().nullish(), + dueDate: nullableText, + branchName: nullableText, + createdAt: nullableText, + updatedAt: nullableText, + state: named.nullish(), + team: named.extend({ key: nullableText }).nullish(), + project: named.nullish(), + cycle: named.nullish(), + assignee: z + .object({ id: nullableText, displayName: nullableText, avatarUrl: nullableText }) + .nullish(), + labels: z + .object({ + nodes: z.array(named).max(50).optional(), + pageInfo: z + .object({ + hasNextPage: z.boolean().optional(), + endCursor: nullableText + }) + .optional() + }) + .nullish() +}) +const envelope = z.object({ + data: z.object({ + issues: z.object({ + nodes: z.array(issue), + pageInfo: z.object({ + hasNextPage: z.boolean(), + endCursor: nullableText + }) + }) + }) +}) + +export const LIST_ISSUES_QUERY = `query OrcaLinearListIssues( + $first: Int!, $after: String, $filter: IssueFilter, + $orderBy: PaginationOrderBy, $includeArchived: Boolean +) { issues(first: $first, after: $after, filter: $filter, + orderBy: $orderBy, includeArchived: $includeArchived) { + nodes { ${ISSUE_FIELDS} } pageInfo { hasNextPage endCursor } +} }` + +export async function acquireIssueListPage( + options: LinearClientOptions, + variables: Record & { first: number }, + signal: AbortSignal +): Promise['data']['issues']> { + const { apiKey, accessToken, apiUrl, headers: suppliedHeaders, ...init } = options + const headers = new Headers({ + 'Content-Type': 'application/json', + Authorization: accessToken + ? accessToken.startsWith('Bearer ') + ? accessToken + : `Bearer ${accessToken}` + : (apiKey ?? '') + }) + new Headers(suppliedHeaders).forEach((value, name) => headers.set(name, value)) + const response = await fetch(apiUrl ?? 'https://api.linear.app/graphql', { + ...init, + method: 'POST', + headers, + body: JSON.stringify({ query: LIST_ISSUES_QUERY, variables }), + signal + }) + if (!(response instanceof Response)) { + throw linearError('linear_list_invalid_response', 'Linear requires a streaming response.') + } + if (!response.ok) { + await response.body?.cancel().catch(() => undefined) + const code = + response.status === 401 + ? 'linear_auth_expired' + : response.status === 403 + ? 'linear_permission_denied' + : response.status === 429 + ? 'linear_rate_limited' + : 'linear_network_error' + const retryAfter = Number(response.headers.get('retry-after')) + throw linearError( + code, + `Linear provider request failed (HTTP ${response.status}).`, + Number.isFinite(retryAfter) && retryAfter >= 0 && retryAfter <= 86_400 + ? { retryAfterSeconds: retryAfter } + : undefined + ) + } + const bytes = await readFetchResponseBytesWithinLimit(response, 4 * 1024 * 1024, signal) + signal.throwIfAborted() + let content: string + try { + content = new TextDecoder('utf-8', { fatal: true }).decode(bytes) + } catch { + throw linearError('linear_list_invalid_response', 'Linear returned invalid UTF-8.') + } + assertJsonTextStructureWithinLimits(content, { structuralTokens: 250_000, nestingDepth: 32 }) + let raw: unknown + try { + raw = JSON.parse(content) + } catch { + throw linearError('linear_list_invalid_response', 'Linear returned invalid JSON.') + } + if ( + raw && + typeof raw === 'object' && + 'errors' in raw && + (!Array.isArray(raw.errors) || raw.errors.length > 0) + ) { + const firstError = Array.isArray(raw.errors) ? raw.errors[0] : undefined + const type = firstError?.extensions?.type ?? firstError?.message + const code = + type === 'AuthenticationError' + ? 'linear_auth_expired' + : type === 'Forbidden' + ? 'linear_permission_denied' + : type === 'Ratelimited' + ? 'linear_rate_limited' + : 'linear_network_error' + throw linearError(code, 'Linear returned a GraphQL error.') + } + const parsed = envelope.safeParse(raw) + if (!parsed.success || parsed.data.data.issues.nodes.length > variables.first) { + throw linearError('linear_list_invalid_response', 'Linear returned an invalid issue page.') + } + return parsed.data.data.issues +} diff --git a/src/main/linear/mcp-issue-list-admission.ts b/src/main/linear/mcp-issue-list-admission.ts new file mode 100644 index 00000000000..c96fc4121a3 --- /dev/null +++ b/src/main/linear/mcp-issue-list-admission.ts @@ -0,0 +1,82 @@ +import { createHash } from 'node:crypto' +import type { LinearMcpIssueListResult } from '../../shared/linear/mcp-issue-list' +import { + stringifyJsonWithinByteLimit, + JsonStringifyByteLimitError +} from '../../shared/node-bounded-json-stringify' +import { linearError } from './issue-context-errors' + +export const LIST_ISSUE_BYTES = 896 * 1024 +export type ListedIssue = LinearMcpIssueListResult['issues'][number] +export class IssueListAdmission { + readonly issues: ListedIssue[] = [] + private bytes = 0 + private readonly identities = new Set() + private readonly cursors = new Set() + private progressBytes = 0 + + get remainingBytes(): number { + return LIST_ISSUE_BYTES - this.bytes + } + + stage( + rows: ListedIssue[], + cursorIdentity?: string + ): { bytes: number; commit: () => void } | null { + let bytes = 0 + const identities = rows.map((row) => digest(JSON.stringify([row.workspace.id, row.id]))) + if ( + new Set(identities).size !== identities.length || + identities.some((key) => this.identities.has(key)) + ) { + throw linearError( + 'linear_list_invalid_response', + 'Linear returned duplicate or conflicting issue identities.' + ) + } + const cursor = cursorIdentity === undefined ? undefined : digest(cursorIdentity) + if (cursor && this.cursors.has(cursor)) { + throw linearError('linear_list_cursor_cycle', 'Linear returned a repeated provider cursor.') + } + const progressBytes = (identities.length + (cursor ? 1 : 0)) * 64 + if (this.progressBytes + progressBytes > 64 * 1024) { + throw linearError( + 'linear_list_metadata_capacity', + 'Linear progress metadata capacity reached; resume from the returned position.' + ) + } + try { + for (const row of rows) { + const wrapper = { result: { issues: [row] } } + bytes += Math.max( + stringifyJsonWithinByteLimit(wrapper, LIST_ISSUE_BYTES).byteLength, + stringifyJsonWithinByteLimit(wrapper, LIST_ISSUE_BYTES, 2).byteLength + 1 + ) + } + } catch (error) { + if (error instanceof JsonStringifyByteLimitError) { + return null + } + throw error + } + if (bytes > this.remainingBytes) { + return null + } + return { + bytes, + commit: () => { + this.issues.push(...rows) + this.bytes += bytes + this.progressBytes += progressBytes + identities.forEach((key) => this.identities.add(key)) + if (cursor) { + this.cursors.add(cursor) + } + } + } + } +} + +function digest(value: string): string { + return createHash('sha256').update(value).digest('hex') +} diff --git a/src/main/linear/mcp-issue-list-lifetime.test.ts b/src/main/linear/mcp-issue-list-lifetime.test.ts new file mode 100644 index 00000000000..556233ff60c --- /dev/null +++ b/src/main/linear/mcp-issue-list-lifetime.test.ts @@ -0,0 +1,94 @@ +import { afterEach, describe, expect, it, vi } from 'vitest' +import { IssueListLifetime } from './mcp-issue-list-lifetime' +import { acquire, release } from './linear-request-concurrency' +import { readFetchResponseBytesWithinLimit } from '../../shared/fetch-response-body' +vi.mock('./linear-token-store', () => ({ clearToken: vi.fn() })) + +function deferred() { + let resolve!: (value: T) => void + const promise = new Promise((r) => { + resolve = r + }) + return { promise, resolve } +} +afterEach(() => vi.useRealTimers()) + +describe('Linear list lease lifetime', () => { + it('removes an aborted queued waiter without admitting it later', async () => { + await Promise.all([acquire(), acquire(), acquire(), acquire()]) + const abort = new AbortController() + const queued = acquire(abort.signal) + const observed = expect(queued).rejects.toBe('cancelled') + abort.abort('cancelled') + await observed + release() + release() + release() + release() + await Promise.all([acquire(), acquire(), acquire(), acquire()]) + release() + release() + release() + release() + }) + it.each(['read', 'cancel'] as const)( + 'holds list reservation and provider lease until stuck %s settles', + async (stuck) => { + vi.useFakeTimers() + const read = deferred>() + const cancel = deferred() + const reader = { + read: () => read.promise, + cancel: () => cancel.promise, + releaseLock: vi.fn() + } + const response = { + headers: new Headers(), + body: { getReader: () => reader } + } as unknown as Response + const owner = new IssueListLifetime(undefined, 10) + const listing = owner.read('fixture', (signal) => + readFetchResponseBytesWithinLimit(response, 1024, signal) + ) + const observed = expect(listing).rejects.toMatchObject({ code: 'linear_timeout' }) + await vi.advanceTimersByTimeAsync(11) + await observed + owner.finish() + const rest = Array.from({ length: 27 }, () => new IssueListLifetime()) + expect(() => new IssueListLifetime()).toThrow('capacity') + await Promise.all([acquire(), acquire(), acquire()]) + let fifthAdmitted = false + const fifth = acquire().then(() => { + fifthAdmitted = true + }) + if (stuck === 'read') { + cancel.resolve() + } else { + read.resolve({ done: true, value: undefined }) + } + await vi.advanceTimersByTimeAsync(1) + expect(fifthAdmitted).toBe(false) + expect(() => new IssueListLifetime()).toThrow('capacity') + read.resolve({ done: true, value: undefined }) + cancel.resolve() + await fifth + const recovered = new IssueListLifetime() + recovered.finish() + rest.forEach((item) => item.finish()) + release() + release() + release() + release() + expect(reader.releaseLock).toHaveBeenCalledOnce() + } + ) + it('holds a completed result until delivery handoff', async () => { + const owners = Array.from({ length: 28 }, () => new IssueListLifetime()) + expect(await owners[0].read('fixture', async () => 'complete')).toBe('complete') + expect(() => new IssueListLifetime()).toThrow('capacity') + owners[0].finish() + const next = new IssueListLifetime() + next.finish() + owners.forEach((owner) => owner.finish()) + }) +}) diff --git a/src/main/linear/mcp-issue-list-lifetime.ts b/src/main/linear/mcp-issue-list-lifetime.ts new file mode 100644 index 00000000000..52239a299cc --- /dev/null +++ b/src/main/linear/mcp-issue-list-lifetime.ts @@ -0,0 +1,90 @@ +import { acquire, release, reserveLinearListing } from './linear-request-concurrency' +import { registerLinearAccountRead } from './linear-account-read-lifetime' +import { LinearAgentAccessError, linearError } from './issue-context-errors' +import { clearToken } from './linear-token-store' + +export class IssueListLifetime { + private readonly releaseReservation: () => void + private pending = 0 + private finished = false + readonly deadline: number + + constructor( + readonly signal?: AbortSignal, + budgetMs = 20_000 + ) { + const reservation = reserveLinearListing() + if (!reservation) { + throw linearError('linear_list_capacity', 'Linear listing capacity is busy; retry later.') + } + this.releaseReservation = reservation + this.deadline = Date.now() + budgetMs + } + + finish(): void { + this.finished = true + this.maybeRelease() + } + + private maybeRelease(): void { + if (this.finished && this.pending === 0) { + this.releaseReservation() + } + } + + async read(workspaceId: string, operation: (signal: AbortSignal) => Promise): Promise { + this.signal?.throwIfAborted() + const remaining = this.deadline - Date.now() + if (remaining <= 0) { + throw linearError('linear_timeout', 'Linear listing deadline reached.') + } + const account = registerLinearAccountRead(workspaceId) + const timeout = new AbortController() + const signal = AbortSignal.any([ + account.signal, + timeout.signal, + ...(this.signal ? [this.signal] : []) + ]) + const timer = setTimeout( + () => timeout.abort(linearError('linear_timeout', 'Linear listing deadline reached.')), + remaining + ) + this.pending++ + const settled = (async () => { + await acquire(signal) + try { + signal.throwIfAborted() + return await operation(signal) + } catch (error) { + if (error instanceof LinearAgentAccessError && error.code === 'linear_auth_expired') { + clearToken(workspaceId) + } + throw error + } finally { + release() + } + })().finally(() => { + clearTimeout(timer) + account.dispose() + this.pending-- + this.maybeRelease() + }) + let abort: (() => void) | undefined + try { + return await Promise.race([ + settled, + new Promise((_, reject) => { + abort = () => reject(signal.reason) + signal.addEventListener('abort', abort, { once: true }) + if (signal.aborted) { + abort() + } + }) + ]) + } finally { + if (abort) { + signal.removeEventListener('abort', abort) + } + } + } +} diff --git a/src/main/linear/mcp-issue-list-page-owned.test.ts b/src/main/linear/mcp-issue-list-page-owned.test.ts new file mode 100644 index 00000000000..8cdbb5a7558 --- /dev/null +++ b/src/main/linear/mcp-issue-list-page-owned.test.ts @@ -0,0 +1,264 @@ +import { afterEach, beforeEach, describe, expect, it, vi } from 'vitest' +import { listMcpIssues } from './mcp-issue-list' +import { decodeIssueListCursor } from './mcp-issue-list-cursor' +import { invalidateLinearAccountReads } from './linear-account-read-lifetime' + +const state = vi.hoisted(() => ({ + workspaces: [ + { + id: 'a', + organizationId: 'a', + organizationName: 'A', + displayName: 'A', + email: null, + credentialRevision: 1 + } + ] +})) +vi.mock('./client', () => ({ + getStatus: () => ({ workspaces: state.workspaces, activeWorkspaceId: 'a' }), + getClients: (id: string) => + state.workspaces + .filter((w) => w.id === id) + .map((workspace) => ({ workspace, apiKey: workspace.id })) +})) +vi.mock('./linear-token-store', () => ({ clearToken: vi.fn() })) + +function row(id: number, size = 1) { + return { + id: String(id), + identifier: `I-${id}`, + title: 'Issue', + url: 'https://linear.app/example', + description: 'x'.repeat(size), + updatedAt: '2026-01-01', + priority: 2 + } +} +function provider(rows: ReturnType[]) { + const calls: { first: number; after?: string }[] = [] + vi.stubGlobal( + 'fetch', + vi.fn(async (_url, options) => { + const { variables } = JSON.parse(options.body) + calls.push(variables) + const offset = Number(variables.after ?? 0) + const nodes = rows.slice(offset, offset + variables.first) + return Response.json({ + data: { + issues: { + nodes, + pageInfo: { + hasNextPage: offset + nodes.length < rows.length, + endCursor: String(offset + nodes.length) + } + } + } + }) + }) + ) + return calls +} + +afterEach(() => { + vi.unstubAllGlobals() + vi.useRealTimers() +}) +beforeEach(() => { + state.workspaces = [ + { + id: 'a', + organizationId: 'a', + organizationName: 'A', + displayName: 'A', + email: null, + credentialRevision: 1 + } + ] +}) + +describe('actual page-owned Linear producer', () => { + it('walks omitted500 without a hidden250 limit and preserves full fields', async () => { + const rows = Array.from({ length: 500 }, (_, i) => row(i)) + rows[8] = row(8, 130 * 1024) + const calls = provider(rows) + const result = await listMcpIssues({ workspaceId: 'a' }) + expect(result.issues).toHaveLength(500) + expect(result.issues[8].description).toHaveLength(130 * 1024) + expect(result.issues[8].priorityLabel).toBe('high') + expect(result.meta.hasMore).toBe(false) + expect(calls).toHaveLength(2) + }) + it('preserves concrete v1 limit1 and exact exhaustion', async () => { + provider([row(0), row(1)]) + const first = await listMcpIssues({ workspaceId: 'a', limit: 1 }) + expect(decodeIssueListCursor(first.meta.nextCursor!)?.cursor).toBe('1') + const next = await listMcpIssues({ cursor: first.meta.nextCursor, limit: 1 }) + expect(next.issues[0].id).toBe('1') + expect(next.meta.hasMore).toBe(false) + expect(next.meta.nextCursor).toBeUndefined() + }) + it('retries whole pages without gaps under heterogeneous byte pressure', async () => { + const rows = Array.from({ length: 80 }, (_, i) => row(i, i % 7 === 0 ? 130 * 1024 : 45 * 1024)) + const calls = provider(rows) + const ids: string[] = [] + let cursor: string | undefined + for (let invocation = 0; invocation < 30; invocation++) { + const result = await listMcpIssues({ workspaceId: 'a', cursor }) + ids.push(...result.issues.map((issue) => issue.id)) + if (!result.meta.hasMore) { + break + } + cursor = result.meta.nextCursor + expect(cursor).toBeTruthy() + } + expect(ids).toEqual(rows.map((issue) => issue.id)) + expect(calls[0].first).toBe(250) + expect(calls[1].after).toBeUndefined() + expect(calls[1].first).toBeLessThan(250) + }) + it('errors on a sole unrepresentable record and does not skip it', async () => { + provider([row(0, 950 * 1024), row(1)]) + await expect(listMcpIssues({ workspaceId: 'a' })).rejects.toMatchObject({ + code: 'linear_list_record_too_large', + data: { retryPosition: { workspaceId: 'a' } } + }) + }) + it('distinguishes unknown acquisition size from record size', async () => { + provider([row(0, 5 * 1024 * 1024)]) + await expect(listMcpIssues({ workspaceId: 'a', limit: 1 })).rejects.toMatchObject({ + code: 'linear_list_acquisition_too_large' + }) + }) + it.each([401, 403, 429, 503])('classifies HTTP %s before oversized body', async (status) => { + vi.stubGlobal( + 'fetch', + vi.fn(async () => new Response('x'.repeat(5 * 1024 * 1024), { status })) + ) + const code = + status === 401 + ? 'linear_auth_expired' + : status === 403 + ? 'linear_permission_denied' + : status === 429 + ? 'linear_rate_limited' + : 'linear_network_error' + await expect(listMcpIssues({ workspaceId: 'a' })).rejects.toMatchObject({ code }) + }) + it('rejects invalid UTF8 before mapping', async () => { + vi.stubGlobal( + 'fetch', + vi.fn(async () => new Response(new Uint8Array([0xff]))) + ) + await expect(listMcpIssues({ workspaceId: 'a' })).rejects.toMatchObject({ + code: 'linear_list_invalid_response' + }) + }) + it('distinguishes empty nonterminal page from exhaustion', async () => { + vi.stubGlobal( + 'fetch', + vi.fn(async () => + Response.json({ + data: { issues: { nodes: [], pageInfo: { hasNextPage: true, endCursor: 'next' } } } + }) + ) + ) + await expect(listMcpIssues({ workspaceId: 'a' })).rejects.toMatchObject({ + code: 'linear_list_empty_page' + }) + provider([]) + expect((await listMcpIssues({ workspaceId: 'a' })).meta.hasMore).toBe(false) + }) + it('rotates all-workspace limit1 calls and explicitly refuses legacy incomplete-all', async () => { + state.workspaces.push({ ...state.workspaces[0], id: 'b', organizationId: 'b' }) + provider([row(0)]) + await expect(listMcpIssues({ workspaceId: 'all', limit: 1 })).rejects.toMatchObject({ + code: 'linear_list_concrete_workspace_required' + }) + const first = await listMcpIssues({ + workspaceId: 'all', + limit: 1, + pageRecovery: { version: 1 } + }) + const second = await listMcpIssues({ + workspaceId: 'all', + limit: 1, + pageRecovery: { version: 1, continuation: first.meta.pageRecovery!.continuation } + }) + expect(first.issues[0].workspace.id).toBe('a') + expect(second.issues[0].workspace.id).toBe('b') + expect(second.meta.hasMore).toBe(false) + }) + it('rotates an actually timed-out workspace and resumes healthy work next call', async () => { + vi.useFakeTimers() + state.workspaces.push({ ...state.workspaces[0], id: 'b', organizationId: 'b' }) + let settle!: (value: Response) => void + const held = new Promise((resolve) => { + settle = resolve + }) + vi.stubGlobal( + 'fetch', + vi.fn(async (_url, options) => { + if (options.headers.get('Authorization') === 'a') { + return held + } + return Response.json({ + data: { issues: { nodes: [row(1)], pageInfo: { hasNextPage: false } } } + }) + }) + ) + const pending = listMcpIssues({ + workspaceId: 'all', + limit: 1, + pageRecovery: { version: 1 } + }).catch((error) => error) + await vi.advanceTimersByTimeAsync(20_001) + const failure = await pending + expect(failure.code).toBe('linear_timeout') + const vector = JSON.parse( + Buffer.from(failure.data.pageRecovery.continuation, 'base64url').toString() + ) + expect(vector.nextWorkspaceIndex).toBe(1) + expect(vector.workspaces[0]).not.toHaveProperty('after') + const next = await listMcpIssues({ + workspaceId: 'all', + limit: 1, + pageRecovery: { version: 1, continuation: failure.data.pageRecovery.continuation } + }) + expect(next.issues[0].workspace.id).toBe('b') + settle( + Response.json({ data: { issues: { nodes: [row(0)], pageInfo: { hasNextPage: false } } } }) + ) + await vi.advanceTimersByTimeAsync(1) + }) + it.each([ + {}, + { data: { issues: null } }, + { data: { issues: { nodes: [], pageInfo: null } } }, + { data: { issues: { nodes: [row(0), row(1)], pageInfo: { hasNextPage: false } } } }, + { data: { issues: { nodes: [{ ...row(0), title: 3 }], pageInfo: { hasNextPage: false } } } } + ])('rejects malformed or over-count provider pages without progress', async (body) => { + vi.stubGlobal( + 'fetch', + vi.fn(async () => Response.json(body)) + ) + await expect(listMcpIssues({ workspaceId: 'a', limit: 1 })).rejects.toMatchObject({ + code: 'linear_list_invalid_response', + data: { retryPosition: { workspaceId: 'a' } } + }) + }) + it('suppresses a page invalidated during acquisition', async () => { + vi.stubGlobal( + 'fetch', + vi.fn(async () => { + invalidateLinearAccountReads('a') + return Response.json({ + data: { issues: { nodes: [row(0)], pageInfo: { hasNextPage: false } } } + }) + }) + ) + await expect(listMcpIssues({ workspaceId: 'a' })).rejects.toMatchObject({ + code: 'linear_network_error' + }) + }) +}) diff --git a/src/main/linear/mcp-issue-list-pages.ts b/src/main/linear/mcp-issue-list-pages.ts new file mode 100644 index 00000000000..89047852c41 --- /dev/null +++ b/src/main/linear/mcp-issue-list-pages.ts @@ -0,0 +1,280 @@ +import type { + LinearMcpIssueListRequest, + LinearMcpIssueListResult +} from '../../shared/linear/agent-access' +import { FetchResponseBodyTooLargeError } from '../../shared/fetch-response-body' +import { JsonTextStructureCapacityError } from '../../shared/json-text-structure-limit' +import { getClients, getStatus } from './client' +import { LinearAgentAccessError, linearError } from './issue-context-errors' +import { mapIssue } from './issue-context-raw' +import { encodeIssueListCursor } from './mcp-issue-list-cursor' +import { buildIssueFilter } from './mcp-issue-list-filter' +import { acquireIssueListPage } from './mcp-issue-list-acquisition' +import { IssueListAdmission } from './mcp-issue-list-admission' +import type { IssueListLifetime } from './mcp-issue-list-lifetime' +import { + boundedListJson, + encodePageRecovery, + LIST_CURSOR_BYTES, + type IssueListRecoveryVector +} from './mcp-issue-list-recovery' + +export async function readIssueListPages( + request: LinearMcpIssueListRequest, + state: IssueListRecoveryVector, + owner: IssueListLifetime +): Promise { + const admission = new IssueListAdmission() + const failures: LinearMcpIssueListResult['meta']['workspaceErrors'] = [] + const failed = new Set() + const limit = request.limit === undefined ? null : Math.max(1, Math.floor(request.limit)) + const filter = buildIssueFilter(request) + let attempts = 0 + let average = 0 + let stopReason: string | undefined + let omittedWorkspaceErrors = 0 + while (state.workspaces.some((w) => !w.done && !failed.has(w.id))) { + owner.signal?.throwIfAborted() + if (limit !== null && admission.issues.length >= limit) { + stopReason = 'row_limit' + break + } + if (Date.now() >= owner.deadline || attempts >= 200) { + stopReason = 'budget' + break + } + const index = state.nextWorkspaceIndex + const position = state.workspaces[index] + if (position.done || failed.has(position.id)) { + state.nextWorkspaceIndex = (index + 1) % state.workspaces.length + continue + } + const remaining = limit === null ? 250 : limit - admission.issues.length + let first = Math.min( + 250, + remaining, + average ? Math.max(1, Math.floor(admission.remainingBytes / average)) : 250 + ) + try { + let committed = false + while (!committed) { + if (Date.now() >= owner.deadline || attempts >= 200) { + stopReason = 'budget' + break + } + attempts++ + const entry = getClients(position.id)[0] + if (!entry || (entry.workspace.credentialRevision ?? 0) !== position.credentialRevision) { + throw linearError( + 'linear_list_stale_recovery', + 'Linear account changed during listing; restart and reconcile.' + ) + } + const page = await owner + .read(position.id, async (signal) => { + return acquireIssueListPage( + { apiKey: entry.apiKey }, + { + first, + after: position.after, + filter, + orderBy: request.orderBy ?? 'updatedAt', + includeArchived: request.includeArchived ?? false + }, + signal + ) + }) + .catch((error: unknown) => { + if ( + error instanceof FetchResponseBodyTooLargeError || + error instanceof JsonTextStructureCapacityError + ) { + return null + } + throw error + }) + if (!page) { + if (first > 1) { + first = Math.max(1, Math.floor(first / 2)) + continue + } + throw linearError( + 'linear_list_acquisition_too_large', + 'Linear acquisition exceeds capacity; record size is unknown.' + ) + } + const more = page.pageInfo.hasNextPage + const after = page.pageInfo.endCursor ?? undefined + if (more && page.nodes.length === 0) { + throw linearError('linear_list_empty_page', 'Linear returned an empty nonterminal page.') + } + if (more && !after) { + throw linearError( + 'linear_list_invalid_response', + 'Linear omitted a required provider cursor.' + ) + } + if (more && after === position.after) { + throw linearError( + 'linear_list_cursor_cycle', + 'Linear repeated the current provider cursor.' + ) + } + if (after && Buffer.byteLength(after) > LIST_CURSOR_BYTES) { + throw linearError( + 'linear_list_metadata_capacity', + 'Linear provider cursor exceeds capacity.' + ) + } + const rows = page.nodes.map((raw) => ({ + ...mapIssue(raw), + workspace: { id: entry.workspace.id, name: entry.workspace.organizationName } + })) + const staged = admission.stage( + rows, + more ? JSON.stringify([position.id, after]) : undefined + ) + if (!staged) { + if (first > 1) { + first = Math.max(1, Math.floor(first / 2)) + continue + } + if (!new IssueListAdmission().stage(rows)) { + throw linearError( + 'linear_list_record_too_large', + 'Linear record exceeds listing capacity; details are not complete.', + { detailsComplete: false } + ) + } + stopReason = 'bytes' + break + } + const next = { + ...state, + nextWorkspaceIndex: (index + 1) % state.workspaces.length, + workspaces: state.workspaces.map((w, i) => + i === index ? { ...w, after: more ? after : undefined, done: !more } : w + ) + } + encodePageRecovery(next) + if (more && encodeIssueListCursor(position.id, after!).length > 4096) { + throw linearError( + 'linear_list_metadata_capacity', + 'Linear continuation exceeds capacity.' + ) + } + owner.signal?.throwIfAborted() + const current = getStatus().workspaces?.find((w) => w.id === position.id) + if (!current || (current.credentialRevision ?? 0) !== position.credentialRevision) { + throw linearError( + 'linear_list_stale_recovery', + 'Linear account changed before page commit.' + ) + } + staged.commit() + if (rows.length) { + average = staged.bytes / rows.length + } + Object.assign(state, next) + committed = true + } + if (stopReason) { + break + } + } catch (error) { + owner.signal?.throwIfAborted() + const failure = + error instanceof LinearAgentAccessError + ? error + : linearError( + 'linear_network_error', + 'Linear listing failed; retry from the returned position.' + ) + const item = { + workspace: { id: position.id, name: position.id }, + code: failure.code, + message: failure.message, + data: { + retryPosition: { + workspaceId: position.id, + ...(position.after + ? { cursor: encodeIssueListCursor(position.id, position.after) } + : {}) + }, + detailsComplete: false + } + } + try { + boundedListJson([...failures, item], 32 * 1024) + failures.push(item) + } catch { + omittedWorkspaceErrors++ + } + failed.add(position.id) + state.nextWorkspaceIndex = (index + 1) % state.workspaces.length + } + } + const hasMore = state.workspaces.some((w) => !w.done) + const pageRecovery = request.pageRecovery + ? { + version: 1 as const, + continuation: encodePageRecovery(state), + ordering: 'admitted_batch' as const, + consistency: 'best_effort' as const, + ...(stopReason ? { stopReason } : {}) + } + : undefined + if (request.workspaceId === 'all' && hasMore && !pageRecovery) { + throw linearError( + 'linear_list_concrete_workspace_required', + 'Incomplete all-workspace listing requires concrete workspace restart and reconciliation.' + ) + } + if (admission.issues.length === 0 && hasMore) { + const failure = failures[0] + throw linearError( + failure?.code ?? 'linear_timeout', + failure?.message ?? 'Linear listing stopped before a page was admitted.', + { + ...(pageRecovery + ? { pageRecovery } + : { + retryPosition: { + workspaceId: state.workspaces[0].id, + ...(state.workspaces[0].after + ? { + cursor: encodeIssueListCursor( + state.workspaces[0].id, + state.workspaces[0].after + ) + } + : {}) + } + }), + detailsComplete: false + } + ) + } + admission.issues.sort((a, b) => + (b[request.orderBy ?? 'updatedAt'] ?? '').localeCompare(a[request.orderBy ?? 'updatedAt'] ?? '') + ) + const concrete = request.workspaceId !== 'all' ? state.workspaces[0] : undefined + return { + issues: admission.issues, + truncated: hasMore, + meta: { + limit, + returned: admission.issues.length, + hasMore, + ...(concrete && hasMore && concrete.after + ? { nextCursor: encodeIssueListCursor(concrete.id, concrete.after) } + : {}), + ...(pageRecovery ? { pageRecovery } : {}), + orderBy: request.orderBy ?? 'updatedAt', + workspaceId: concrete?.id ?? 'all', + partial: failures.length + omittedWorkspaceErrors > 0, + workspaceErrors: failures, + ...(omittedWorkspaceErrors ? { omittedWorkspaceErrors } : {}) + } + } +} diff --git a/src/main/linear/mcp-issue-list-recovery.ts b/src/main/linear/mcp-issue-list-recovery.ts new file mode 100644 index 00000000000..7032193140e --- /dev/null +++ b/src/main/linear/mcp-issue-list-recovery.ts @@ -0,0 +1,166 @@ +import { createHash } from 'node:crypto' +import { z } from 'zod' +import type { LinearMcpIssueListRequest } from '../../shared/linear/mcp-issue-list' +import type { LinearWorkspace } from '../../shared/linear/workspace-types' +import { stringifyJsonWithinByteLimit } from '../../shared/node-bounded-json-stringify' +import { linearError } from './issue-context-errors' +import { buildIssueFilter } from './mcp-issue-list-filter' + +export const LIST_CONTEXT_BYTES = 64 * 1024 +export const LIST_CURSOR_BYTES = 2048 +const position = z + .object({ + id: z.string().min(1), + credentialRevision: z.number().int().nonnegative(), + after: z.string().min(1).optional(), + done: z.boolean() + }) + .strict() +const vector = z + .object({ + version: z.literal(1), + queryHash: z.string().length(64), + rosterHash: z.string().length(64), + nextWorkspaceIndex: z.number().int().nonnegative(), + workspaces: z.array(position).min(1) + }) + .strict() +export type IssueListRecoveryVector = z.infer + +export function boundedListJson(value: unknown, maxBytes = LIST_CONTEXT_BYTES): string { + try { + return stringifyJsonWithinByteLimit(value, maxBytes).serialized + } catch { + throw linearError( + 'linear_list_metadata_capacity', + 'Linear list metadata exceeds capacity; restart with a concrete workspace.' + ) + } +} + +function hash(value: unknown): string { + return createHash('sha256').update(boundedListJson(value)).digest('hex') +} + +function canonical(value: unknown): unknown { + if (Array.isArray(value)) { + return value.map(canonical) + } + if (value && typeof value === 'object') { + return Object.fromEntries( + Object.entries(value) + .sort(([a], [b]) => (a < b ? -1 : a > b ? 1 : 0)) + .map(([key, item]) => [key, canonical(item)]) + ) + } + return value +} + +export function encodePageRecovery(state: IssueListRecoveryVector): string { + const encoded = Buffer.from(boundedListJson(state)).toString('base64url') + if (Buffer.byteLength(encoded) > LIST_CONTEXT_BYTES) { + throw linearError( + 'linear_list_metadata_capacity', + 'Linear page recovery exceeds capacity; use concrete workspace recovery.' + ) + } + return encoded +} + +export function createPageRecovery( + request: LinearMcpIssueListRequest, + workspaces: LinearWorkspace[] +): IssueListRecoveryVector { + const queryHash = hash( + canonical({ + filter: buildIssueFilter(request), + orderBy: request.orderBy ?? 'updatedAt', + includeArchived: request.includeArchived ?? false + }) + ) + const roster = workspaces + .map(({ id, credentialRevision }) => ({ id, credentialRevision: credentialRevision ?? 0 })) + .sort((a, b) => (a.id < b.id ? -1 : a.id > b.id ? 1 : 0)) + const rosterHash = hash(roster) + const initial: IssueListRecoveryVector = { + version: 1, + queryHash, + rosterHash, + nextWorkspaceIndex: 0, + workspaces: roster.map((workspace) => ({ ...workspace, done: false })) + } + encodePageRecovery(initial) + const encoded = request.pageRecovery?.continuation + if (!encoded) { + return initial + } + if (Buffer.byteLength(encoded) > LIST_CONTEXT_BYTES || !/^[A-Za-z0-9_-]+$/.test(encoded)) { + throw linearError( + 'linear_list_stale_recovery', + 'Invalid Linear recovery; restart with a concrete workspace.' + ) + } + let parsed: IssueListRecoveryVector + try { + parsed = vector.parse( + JSON.parse( + new TextDecoder('utf-8', { fatal: true }).decode(Buffer.from(encoded, 'base64url')) + ) + ) + } catch { + throw linearError( + 'linear_list_stale_recovery', + 'Invalid Linear recovery; restart with a concrete workspace.' + ) + } + if ( + parsed.queryHash !== queryHash || + parsed.rosterHash !== rosterHash || + parsed.nextWorkspaceIndex >= roster.length || + parsed.workspaces.length !== roster.length || + parsed.workspaces.some( + (entry, index) => + entry.id !== roster[index].id || + entry.credentialRevision !== roster[index].credentialRevision || + (entry.after !== undefined && Buffer.byteLength(entry.after) > LIST_CURSOR_BYTES) + ) + ) { + throw linearError( + 'linear_list_stale_recovery', + 'Linear query or accounts changed; restart and reconcile by workspace and issue ID.' + ) + } + return parsed +} + +export function decodeDeliveryPageRecovery( + value: unknown +): { version: 1; continuation?: string } | undefined { + if (!value || typeof value !== 'object' || !('version' in value) || value.version !== 1) { + return undefined + } + if (!('continuation' in value) || value.continuation === undefined) { + return { version: 1 } + } + const encoded = value.continuation + if ( + typeof encoded !== 'string' || + Buffer.byteLength(encoded) > LIST_CONTEXT_BYTES || + !/^[A-Za-z0-9_-]+$/.test(encoded) + ) { + return undefined + } + try { + const state = vector.parse( + JSON.parse( + new TextDecoder('utf-8', { fatal: true }).decode(Buffer.from(encoded, 'base64url')) + ) + ) + if (state.workspaces.some((w) => w.after && Buffer.byteLength(w.after) > LIST_CURSOR_BYTES)) { + return undefined + } + return { version: 1, continuation: encodePageRecovery(state) } + } catch { + return undefined + } +} diff --git a/src/main/linear/mcp-issue-list.test.ts b/src/main/linear/mcp-issue-list.test.ts index e2334628f62..363f3b732cd 100644 --- a/src/main/linear/mcp-issue-list.test.ts +++ b/src/main/linear/mcp-issue-list.test.ts @@ -1,4 +1,4 @@ -import { beforeEach, describe, expect, it, vi } from 'vitest' +import { afterEach, beforeEach, describe, expect, it, vi } from 'vitest' const rawRequest = vi.fn() const getClients = vi.fn() @@ -21,12 +21,14 @@ const clientEntry = ( request: ReturnType = rawRequest ) => ({ workspace: workspace(id, organizationName), + apiKey: id, client: { client: { rawRequest: request } } }) vi.mock('./linear-request-concurrency', () => ({ acquire, - release + release, + reserveLinearListing: () => () => {} })) vi.mock('./linear-token-store', () => ({ @@ -40,11 +42,21 @@ vi.mock('./client', () => ({ })) describe('MCP-compatible Linear issue listing', () => { + afterEach(() => vi.unstubAllGlobals()) beforeEach(() => { vi.clearAllMocks() const entry = clientEntry('workspace-1', 'Acme') getClients.mockReturnValue([entry]) getStatus.mockReturnValue({ workspaces: [entry.workspace] }) + vi.stubGlobal( + 'fetch', + vi.fn(async (_url, options) => { + const { query, variables } = JSON.parse(options.body) + const selected = getClients(options.headers.get('Authorization'))[0] + const body = await selected.client.client.rawRequest(query, variables) + return body instanceof Response ? body : Response.json(body) + }) + ) }) it('passes rich filters, ordering, archive scope, and cursor to Linear', async () => { @@ -157,7 +169,7 @@ describe('MCP-compatible Linear issue listing', () => { }) }) - it('fans out one bounded provider request per workspace concurrently', async () => { + it('takes one provider page at a time and sorts the admitted batch', async () => { const firstRequest = vi.fn() const secondRequest = vi.fn() let resolveFirst: ((value: unknown) => void) | undefined @@ -182,10 +194,10 @@ describe('MCP-compatible Linear issue listing', () => { }) const { listMcpIssues } = await import('./mcp-issue-list') - const pending = listMcpIssues({ limit: 1, workspaceId: 'all' }) + const pending = listMcpIssues({ limit: 2, workspaceId: 'all', pageRecovery: { version: 1 } }) await vi.waitFor(() => { expect(firstRequest).toHaveBeenCalledTimes(1) - expect(secondRequest).toHaveBeenCalledTimes(1) + expect(secondRequest).not.toHaveBeenCalled() }) resolveFirst?.({ data: { @@ -195,6 +207,7 @@ describe('MCP-compatible Linear issue listing', () => { } } }) + await vi.waitFor(() => expect(secondRequest).toHaveBeenCalledTimes(1)) resolveSecond?.({ data: { issues: { @@ -205,10 +218,10 @@ describe('MCP-compatible Linear issue listing', () => { }) const result = await pending - expect(result.issues.map((issue) => issue.identifier)).toEqual(['OPS-1']) - expect(result.meta).toMatchObject({ returned: 1, hasMore: true, partial: false }) + expect(result.issues.map((issue) => issue.identifier)).toEqual(['OPS-1', 'ENG-1']) + expect(result.meta).toMatchObject({ returned: 2, hasMore: false, partial: false }) expect(result.meta.nextCursor).toBeUndefined() - expect(firstRequest.mock.calls[0]?.[1]).toMatchObject({ first: 1 }) + expect(firstRequest.mock.calls[0]?.[1]).toMatchObject({ first: 2 }) expect(secondRequest.mock.calls[0]?.[1]).toMatchObject({ first: 1 }) }) @@ -221,7 +234,7 @@ describe('MCP-compatible Linear issue listing', () => { } } }) - const failedRequest = vi.fn().mockRejectedValue(new Error('429 rate limit exceeded')) + const failedRequest = vi.fn().mockResolvedValue(new Response('rate limited', { status: 429 })) const healthy = clientEntry('workspace-1', 'Acme', healthyRequest) const failed = clientEntry('workspace-2', 'Beta', failedRequest) getStatus.mockReturnValue({ workspaces: [healthy.workspace, failed.workspace] }) @@ -236,15 +249,15 @@ describe('MCP-compatible Linear issue listing', () => { }) const { listMcpIssues } = await import('./mcp-issue-list') - const result = await listMcpIssues({ workspaceId: 'all' }) + const result = await listMcpIssues({ workspaceId: 'all', pageRecovery: { version: 1 } }) expect(result.issues.map((issue) => issue.identifier)).toEqual(['ENG-1']) expect(result.meta).toMatchObject({ partial: true, returned: 1 }) - expect(result.meta.workspaceErrors).toEqual([ + expect(result.meta.workspaceErrors).toMatchObject([ { - workspace: { id: 'workspace-2', name: 'Beta' }, + workspace: { id: 'workspace-2', name: 'workspace-2' }, code: 'linear_rate_limited', - message: '429 rate limit exceeded' + message: 'Linear provider request failed (HTTP 429).' } ]) }) @@ -272,7 +285,11 @@ describe('MCP-compatible Linear issue listing', () => { await expect(listMcpIssues({ cursor: 'next' })).rejects.toMatchObject({ code: 'linear_invalid_workspace' }) - const result = await listMcpIssues({ workspaceId: 'all', limit: 1 }) + const result = await listMcpIssues({ + workspaceId: 'all', + limit: 1, + pageRecovery: { version: 1 } + }) expect(result.meta).toMatchObject({ hasMore: true, workspaceId: 'all' }) expect(result.meta.nextCursor).toBeUndefined() diff --git a/src/main/linear/mcp-issue-list.ts b/src/main/linear/mcp-issue-list.ts index 30811087d1d..0d6d17c49d8 100644 --- a/src/main/linear/mcp-issue-list.ts +++ b/src/main/linear/mcp-issue-list.ts @@ -2,231 +2,69 @@ import type { LinearMcpIssueListRequest, LinearMcpIssueListResult } from '../../shared/linear/agent-access' -import { getClients, getStatus, type LinearClientForWorkspace } from './client' -import { withLinearRead } from './issue-context-client' +import { getStatus } from './client' import { linearError } from './issue-context-errors' -import { - getFanoutClientEntries, - workspaceFailure, - type WorkspaceReadFailure -} from './issue-context-fanout' -import { ISSUE_FIELDS, mapIssue, type RawIssue } from './issue-context-raw' import { resolveWorkspaceSelector } from './issue-context-workspaces' -import { encodeIssueListCursor, resolveIssueListCursor } from './mcp-issue-list-cursor' -import { buildIssueFilter } from './mcp-issue-list-filter' +import { resolveIssueListCursor } from './mcp-issue-list-cursor' +import { IssueListLifetime } from './mcp-issue-list-lifetime' +import { boundedListJson, createPageRecovery, LIST_CURSOR_BYTES } from './mcp-issue-list-recovery' +import { readIssueListPages } from './mcp-issue-list-pages' -// Why: `--limit` is opt-in, so the default read walks every page instead of quietly cutting -// the answer off. Linear caps `first` at 250. The budget and page ceiling are backstops, not -// caps: the CLI abandons an RPC at 60s, so a walk that would outlive it has to stop early and -// say `truncated` with a continuation cursor rather than fail the whole command. -const LIST_ISSUES_PAGE_SIZE = 250 -const LIST_ISSUES_MAX_PAGES = 200 -const LIST_ISSUES_READ_BUDGET_MS = 20_000 - -type RawListIssuesResponse = { - issues?: { - nodes?: RawIssue[] - pageInfo?: { hasNextPage?: boolean; endCursor?: string | null } - } | null -} - -type WorkspaceIssuePage = { - issues: LinearMcpIssueListResult['issues'] - hasMore: boolean - nextCursor?: string -} - -// limit null means unbounded; the deadline still applies. -type IssueListReadBudget = { limit: number | null; deadline: number } - -const LIST_ISSUES_QUERY = ` - query OrcaLinearListIssues( - $first: Int! - $after: String - $filter: IssueFilter - $orderBy: PaginationOrderBy - $includeArchived: Boolean - ) { - issues( - first: $first - after: $after - filter: $filter - orderBy: $orderBy - includeArchived: $includeArchived - ) { - nodes { ${ISSUE_FIELDS} } - pageInfo { hasNextPage endCursor } - } - } -` +type ListOptions = { signal?: AbortSignal; retainUntilDelivery?: (release: () => void) => void } export async function listMcpIssues( - request: LinearMcpIssueListRequest + request: LinearMcpIssueListRequest, + options: ListOptions = {} ): Promise { + boundedListJson(request) + if ( + request.pageRecovery && + (request.workspaceId !== 'all' || request.cursor || request.pageRecovery.version !== 1) + ) { + throw linearError( + 'linear_invalid_workspace', + 'Page recovery requires --workspace all and cannot be combined with --cursor.' + ) + } + if (request.cursor && Buffer.byteLength(request.cursor) > 4096) { + throw linearError( + 'linear_list_metadata_capacity', + 'Linear cursor exceeds capacity; restart with a concrete workspace.' + ) + } const pagination = resolveIssueListCursor(request) - const limit = resolveLimit(request.limit) - const orderBy = request.orderBy ?? 'updatedAt' - const { entries, failures: entryFailures } = getIssueListEntries(pagination.workspaceId) - if (entries.length === 0) { - if (entryFailures[0]) { - throw entryFailures[0].error - } - throw linearError('linear_not_connected', 'Linear is not connected.', { - nextSteps: ['Connect Linear from Orca settings, then retry the issue list.'] - }) + if (pagination.linearCursor && Buffer.byteLength(pagination.linearCursor) > LIST_CURSOR_BYTES) { + throw linearError('linear_list_metadata_capacity', 'Linear provider cursor exceeds capacity.') } - const pagedRequest = { - ...request, - cursor: pagination.linearCursor, - workspaceId: pagination.workspaceId + const status = getStatus() + const selected = + pagination.workspaceId === 'all' + ? (status.workspaces ?? []) + : [ + resolveWorkspaceSelector( + { workspaceId: pagination.workspaceId }, + status.workspaces ?? [] + ) ?? + (status.workspaces ?? []).find((w) => w.id === status.activeWorkspaceId) ?? + status.workspaces?.[0] + ].filter((w): w is NonNullable => !!w) + if (!selected.length) { + throw linearError('linear_not_connected', 'Linear is not connected.') } - // One deadline for the whole call, so fanning out over many workspaces cannot multiply it. - const deadline = Date.now() + LIST_ISSUES_READ_BUDGET_MS - const { pages, failures } = await readIssueListWorkspaces( - entries, - pagedRequest, - { limit, deadline }, - orderBy, - entryFailures - ) - const issues = pages.flatMap((page) => page.issues) - let hasMore = pages.some((page) => page.hasMore) - - issues.sort((left, right) => compareIssues(left, right, orderBy)) - if (limit !== null && issues.length > limit) { - hasMore = true - issues.length = limit + boundedListJson(selected) + const state = createPageRecovery(request, selected) + if (!request.pageRecovery && pagination.linearCursor) { + state.workspaces[0].after = pagination.linearCursor } - const workspaceId = pagination.workspaceId === 'all' ? 'all' : entries[0].workspace.id - return { - issues, - truncated: hasMore, - meta: { - limit, - returned: issues.length, - hasMore, - ...(hasMore && workspaceId !== 'all' && pages.length === 1 && pages[0].nextCursor - ? { nextCursor: encodeIssueListCursor(workspaceId, pages[0].nextCursor) } - : {}), - orderBy, - workspaceId, - partial: failures.length > 0, - workspaceErrors: failures.map(({ workspace, code, message }) => ({ - workspace, - code, - message - })) + const owner = new IssueListLifetime(options.signal) + if (options.retainUntilDelivery) { + options.retainUntilDelivery(() => owner.finish()) + } + try { + return await readIssueListPages(request, state, owner) + } finally { + if (!options.retainUntilDelivery) { + owner.finish() } } } - -function getIssueListEntries(workspaceId?: (string & {}) | 'all'): { - entries: LinearClientForWorkspace[] - failures: WorkspaceReadFailure[] -} { - if (workspaceId === 'all') { - return getFanoutClientEntries() - } - if (workspaceId) { - resolveWorkspaceSelector({ workspaceId }, getStatus().workspaces ?? []) - } - return { entries: getClients(workspaceId), failures: [] } -} - -async function readIssueListWorkspaces( - entries: LinearClientForWorkspace[], - request: LinearMcpIssueListRequest, - budget: IssueListReadBudget, - orderBy: 'createdAt' | 'updatedAt', - initialFailures: WorkspaceReadFailure[] -): Promise<{ pages: WorkspaceIssuePage[]; failures: WorkspaceReadFailure[] }> { - if (request.workspaceId !== 'all') { - return { - pages: [await readIssueListWorkspace(entries[0], request, budget, orderBy)], - failures: [] - } - } - - const settled = await Promise.allSettled( - entries.map((entry) => readIssueListWorkspace(entry, request, budget, orderBy)) - ) - const pages: WorkspaceIssuePage[] = [] - const failures = [...initialFailures] - for (let index = 0; index < settled.length; index += 1) { - const result = settled[index] - if (result.status === 'fulfilled') { - pages.push(result.value) - continue - } - failures.push(workspaceFailure(entries[index].workspace, result.reason)) - } - if (pages.length === 0 && failures.length === entries.length + initialFailures.length) { - throw failures[0].error - } - return { pages, failures } -} - -async function readIssueListWorkspace( - entry: LinearClientForWorkspace, - request: LinearMcpIssueListRequest, - { limit, deadline }: IssueListReadBudget, - orderBy: 'createdAt' | 'updatedAt' -): Promise { - const filter = buildIssueFilter(request) - const issues: WorkspaceIssuePage['issues'] = [] - let after = request.cursor - let hasMore = false - let nextCursor: string | undefined - for (let page = 0; page < LIST_ISSUES_MAX_PAGES; page += 1) { - const first = - limit === null - ? LIST_ISSUES_PAGE_SIZE - : Math.min(limit - issues.length, LIST_ISSUES_PAGE_SIZE) - // Each page takes its own concurrency slot so a long walk cannot starve other reads. - const connection = await withLinearRead(entry, async () => { - const raw = await entry.client.client.rawRequest< - RawListIssuesResponse, - Record - >(LIST_ISSUES_QUERY, { - first, - after, - filter, - orderBy, - includeArchived: request.includeArchived ?? false - }) - return raw.data?.issues - }) - for (const issue of connection?.nodes ?? []) { - issues.push({ - ...mapIssue(issue), - workspace: { id: entry.workspace.id, name: entry.workspace.organizationName } - }) - } - hasMore = connection?.pageInfo?.hasNextPage === true - nextCursor = connection?.pageInfo?.endCursor ?? undefined - if (!hasMore || !nextCursor) { - break - } - if (limit !== null && issues.length >= limit) { - break - } - if (Date.now() >= deadline) { - break - } - after = nextCursor - } - return { issues, hasMore, nextCursor } -} - -// null means unbounded: read until Linear stops handing out pages. -function resolveLimit(limit: number | undefined): number | null { - return limit === undefined ? null : Math.max(1, Math.floor(limit)) -} - -function compareIssues( - left: LinearMcpIssueListResult['issues'][number], - right: LinearMcpIssueListResult['issues'][number], - orderBy: 'createdAt' | 'updatedAt' -): number { - return (right[orderBy] ?? '').localeCompare(left[orderBy] ?? '') -} diff --git a/src/main/runtime/rpc/core.ts b/src/main/runtime/rpc/core.ts index 702ea1b3aaa..1d4d70b6dd5 100644 --- a/src/main/runtime/rpc/core.ts +++ b/src/main/runtime/rpc/core.ts @@ -64,6 +64,7 @@ export type LegacyCoordinatorAuthorityProof = Readonly<{ export type RpcContext = { runtime: OrcaRuntimeService // Why: lets long-poll handlers release immediately on client disconnect instead of running down timeoutMs. See design doc §3.1. + retainUntilDelivery?: (release: () => void) => void signal?: AbortSignal // Why: per-WebSocket key so the server reaps a closing socket's subscriptions without touching sibling sockets sharing the deviceToken. connectionId?: string diff --git a/src/main/runtime/rpc/dispatcher-stream-options.ts b/src/main/runtime/rpc/dispatcher-stream-options.ts index e3151c66b0e..3c1b26f2430 100644 --- a/src/main/runtime/rpc/dispatcher-stream-options.ts +++ b/src/main/runtime/rpc/dispatcher-stream-options.ts @@ -5,6 +5,7 @@ import type { PairingRpcContext } from './core' export type RpcDispatchStreamingOptions = { authenticatedCallerFingerprint?: string connectionId?: string + retainUntilDelivery?: (release: () => void) => void signal?: AbortSignal clientId?: string pairedDeviceId?: string diff --git a/src/main/runtime/rpc/dispatcher.ts b/src/main/runtime/rpc/dispatcher.ts index 73cfa596dd5..6d5986364f0 100644 --- a/src/main/runtime/rpc/dispatcher.ts +++ b/src/main/runtime/rpc/dispatcher.ts @@ -1,3 +1,4 @@ +import { boundLinearListReply, rejectOversizedLinearListRequest } from './linear-list-reply-budget' import { buildRegistry, isStreamingMethod, @@ -24,7 +25,11 @@ import { parseRpcRequestParams } from './dispatcher-request-parsing' import { RpcStreamingDispatcher } from './rpc-streaming-dispatcher' import { invokeDispatcherUnaryMethod } from './dispatcher-unary-method-invocation' -export type DispatcherOptions = { runtime: OrcaRuntimeService; methods?: readonly RpcAnyMethod[] } +export type DispatcherOptions = { + runtime: OrcaRuntimeService + methods?: readonly RpcAnyMethod[] + linearListDelivery?: RpcDispatchStreamingOptions +} type DispatchCallOptions = RpcDispatchStreamingOptions @@ -35,7 +40,10 @@ export class RpcDispatcher { private readonly legacyOrchestration: OrchestrationLegacyCompatibility private readonly streamingDispatcher: RpcStreamingDispatcher - constructor({ runtime, methods = ALL_RPC_METHODS }: DispatcherOptions) { + private readonly linearListDelivery?: RpcDispatchStreamingOptions + + constructor({ runtime, methods = ALL_RPC_METHODS, linearListDelivery }: DispatcherOptions) { + this.linearListDelivery = linearListDelivery this.runtime = runtime this.registry = buildRegistry(methods) this.orchestrationMutations = getOrchestrationMutationExecutor(runtime) @@ -50,6 +58,11 @@ export class RpcDispatcher { } async dispatch(request: RpcRequest, options?: DispatchCallOptions): Promise { + options ??= this.linearListDelivery + const rejected = rejectOversizedLinearListRequest(request) + if (rejected) { + return rejected + } const meta = this.meta() const method = this.registry.get(request.method) if (!method) { @@ -68,7 +81,7 @@ export class RpcDispatcher { const parsedParams = parseRpcRequestParams(request, method, meta) if (parsedParams.error) { - return parsedParams.error + return boundLinearListReply(request, parsedParams.error) } if (isStreamingMethod(method)) { @@ -92,6 +105,7 @@ export class RpcDispatcher { context: { runtime: this.runtime, signal: options?.signal, + retainUntilDelivery: options?.retainUntilDelivery, connectionId: options?.connectionId, requestId: request.id, clientId: options?.clientId, @@ -104,12 +118,12 @@ export class RpcDispatcher { orchestrationMutations: this.orchestrationMutations, legacyOrchestration: this.legacyOrchestration }) - return successResponse(request.id, meta, result) + return boundLinearListReply(request, successResponse(request.id, meta, result)) } catch (error) { if (request.method.startsWith('emulator.')) { emulatorProbeError(`rpc ${request.method}`, error, { params: request.params }) } - return mapDispatcherError(request, meta, error) + return boundLinearListReply(request, mapDispatcherError(request, meta, error)) } } diff --git a/src/main/runtime/rpc/linear-list-reply-budget.ts b/src/main/runtime/rpc/linear-list-reply-budget.ts new file mode 100644 index 00000000000..f473abab605 --- /dev/null +++ b/src/main/runtime/rpc/linear-list-reply-budget.ts @@ -0,0 +1,116 @@ +import { stringifyJsonWithinByteLimit } from '../../../shared/node-bounded-json-stringify' +import type { RpcRequest, RpcResponse } from './core' +import { errorResponse } from './errors' +import { decodeDeliveryPageRecovery } from '../../linear/mcp-issue-list-recovery' + +export function isLinearPageRequest(request: Pick): boolean { + if (request.method === 'linear.mcpListIssues') { + return true + } + return ( + request.method === 'linear.listIssues' && + !!request.params && + typeof request.params === 'object' && + [ + 'team', + 'cycle', + 'label', + 'query', + 'state', + 'cursor', + 'orderBy', + 'project', + 'release', + 'assignee', + 'delegate', + 'parentId', + 'priority', + 'createdAt', + 'updatedAt', + 'includeArchived', + 'pageRecovery' + ].some((key) => key in (request.params as object)) + ) +} + +export function boundLinearListReply(request: RpcRequest, response: RpcResponse): RpcResponse { + if (!isLinearPageRequest(request)) { + return response + } + try { + if (Buffer.byteLength(request.id) > 128) { + throw new Error('correlation capacity') + } + stringifyJsonWithinByteLimit(response._meta, 16 * 1024, 2) + if (!response.ok && Buffer.byteLength(response.error.message) > 512) { + throw new Error('error capacity') + } + const params = request.params as { workspaceId?: string; pageRecovery?: unknown } | undefined + const maxBytes = response.ok + ? 1024 * 1024 + : params?.workspaceId === 'all' && params.pageRecovery + ? 128 * 1024 + : 8192 + stringifyJsonWithinByteLimit(response, maxBytes - 1, 2) + stringifyJsonWithinByteLimit(response, maxBytes) + return response + } catch { + return linearListDeliveryFailure(request) + } +} + +export function linearListDeliveryFailure(request: RpcRequest): RpcResponse { + const params = request.params as + | { workspaceId?: unknown; cursor?: unknown; pageRecovery?: unknown } + | undefined + const id = + typeof request.id === 'string' && Buffer.byteLength(request.id) <= 128 ? request.id : 'unknown' + const workspaceId = + typeof params?.workspaceId === 'string' && Buffer.byteLength(params.workspaceId) <= 2048 + ? params.workspaceId + : undefined + const cursor = + typeof params?.cursor === 'string' && Buffer.byteLength(params.cursor) <= 4096 + ? params.cursor + : undefined + const recovery = { + retryPosition: { ...(workspaceId ? { workspaceId } : {}), ...(cursor ? { cursor } : {}) }, + restartConcreteWorkspaces: workspaceId === 'all' || !workspaceId, + detailsComplete: false + } + const pageRecovery = decodeDeliveryPageRecovery(params?.pageRecovery) + const failure = errorResponse( + id, + { runtimeId: 'unknown' }, + 'linear_list_metadata_capacity', + 'Linear reply could not be delivered; retry the input position or restart concrete workspaces and reconcile by issue ID.', + { ...recovery, ...(workspaceId === 'all' && pageRecovery ? { pageRecovery } : {}) } + ) + try { + stringifyJsonWithinByteLimit(failure, pageRecovery ? 128 * 1024 - 1 : 8191, 2) + return failure + } catch { + return errorResponse( + id, + { runtimeId: 'unknown' }, + 'linear_list_metadata_capacity', + 'Linear reply could not be delivered; restart concrete workspaces and reconcile by issue ID.', + { restartConcreteWorkspaces: true, detailsComplete: false } + ) + } +} + +export function rejectOversizedLinearListRequest(request: RpcRequest): RpcResponse | undefined { + if (!isLinearPageRequest(request)) { + return undefined + } + try { + if (typeof request.id !== 'string' || Buffer.byteLength(request.id) > 128) { + throw new Error('correlation capacity') + } + stringifyJsonWithinByteLimit(request.params, 64 * 1024) + return undefined + } catch { + return linearListDeliveryFailure(request) + } +} diff --git a/src/main/runtime/rpc/methods/linear-issue-list-method.ts b/src/main/runtime/rpc/methods/linear-issue-list-method.ts index fa69b745377..61d9d8917d8 100644 --- a/src/main/runtime/rpc/methods/linear-issue-list-method.ts +++ b/src/main/runtime/rpc/methods/linear-issue-list-method.ts @@ -22,6 +22,10 @@ const McpListIssues = z query: OptionalString, state: OptionalString, cursor: OptionalString, + pageRecovery: z + .object({ version: z.literal(1), continuation: z.string().max(65536).optional() }) + .strict() + .optional(), orderBy: z.enum(['createdAt', 'updatedAt']).optional(), project: OptionalString, release: OptionalString, @@ -41,9 +45,9 @@ const ListIssues = z.union([McpListIssues, LegacyListIssues]) export const LINEAR_ISSUE_LIST_METHOD = defineMethod({ name: 'linear.listIssues', params: ListIssues, - handler: async (params, { runtime }) => { + handler: async (params, { runtime, signal, retainUntilDelivery }) => { if (isMcpIssueListRequest(params)) { - return runtime.linearMcpIssueList(params) + return runtime.linearMcpIssueList(params, { signal, retainUntilDelivery }) } return runtime.linearListIssues(params?.filter, params?.limit, params?.workspaceId, { attributeFilter: params?.attributeFilter @@ -54,7 +58,8 @@ export const LINEAR_ISSUE_LIST_METHOD = defineMethod({ export const LINEAR_MCP_ISSUE_LIST_METHOD = defineMethod({ name: 'linear.mcpListIssues', params: McpListIssues, - handler: async (params, { runtime }) => runtime.linearMcpIssueList(params) + handler: async (params, { runtime, signal, retainUntilDelivery }) => + runtime.linearMcpIssueList(params, { signal, retainUntilDelivery }) }) const MCP_ISSUE_LIST_KEYS = [ @@ -73,7 +78,8 @@ const MCP_ISSUE_LIST_KEYS = [ 'priority', 'createdAt', 'updatedAt', - 'includeArchived' + 'includeArchived', + 'pageRecovery' ] as const function isMcpIssueListRequest( diff --git a/src/main/runtime/rpc/methods/linear.test.ts b/src/main/runtime/rpc/methods/linear.test.ts index 1e9b430fa8a..667ffea16e6 100644 --- a/src/main/runtime/rpc/methods/linear.test.ts +++ b/src/main/runtime/rpc/methods/linear.test.ts @@ -130,17 +130,23 @@ describe('linear RPC methods', () => { expect(runtime.linearListIssues).toHaveBeenCalledWith(undefined, 5, 'workspace-1', { attributeFilter: undefined }) - expect(runtime.linearMcpIssueList).toHaveBeenCalledWith({ - team: 'ENG', - assignee: 'me', - cursor: 'next', - orderBy: 'updatedAt', - workspaceId: 'workspace-1' - }) - expect(runtime.linearMcpIssueList).toHaveBeenCalledWith({ - limit: 5, - workspaceId: 'workspace-1' - }) + expect(runtime.linearMcpIssueList).toHaveBeenCalledWith( + { + team: 'ENG', + assignee: 'me', + cursor: 'next', + orderBy: 'updatedAt', + workspaceId: 'workspace-1' + }, + { signal: undefined, retainUntilDelivery: undefined } + ) + expect(runtime.linearMcpIssueList).toHaveBeenCalledWith( + { + limit: 5, + workspaceId: 'workspace-1' + }, + { signal: undefined, retainUntilDelivery: undefined } + ) expect(runtime.linearGetIssue).toHaveBeenCalledWith('issue-3', 'workspace-1') expect(runtime.linearCreateIssue).toHaveBeenCalledWith( 'team-1', diff --git a/src/main/runtime/rpc/rpc-streaming-dispatcher.ts b/src/main/runtime/rpc/rpc-streaming-dispatcher.ts index 6eddcb748be..160076644f9 100644 --- a/src/main/runtime/rpc/rpc-streaming-dispatcher.ts +++ b/src/main/runtime/rpc/rpc-streaming-dispatcher.ts @@ -1,3 +1,4 @@ +import { boundLinearListReply, rejectOversizedLinearListRequest } from './linear-list-reply-budget' import { isStreamingMethod, type RpcEnvelopeMeta, type RpcRegistry, type RpcRequest } from './core' import { errorResponse, successResponse } from './errors' @@ -35,6 +36,11 @@ export class RpcStreamingDispatcher { ): Promise { const { runtime, registry, orchestrationMutations, legacyOrchestration, meta } = this.dependencies + const rejected = rejectOversizedLinearListRequest(request) + if (rejected) { + reply(JSON.stringify(rejected)) + return + } const envelopeMeta = meta() const method = registry.get(request.method) if (!method) { @@ -59,7 +65,7 @@ export class RpcStreamingDispatcher { const parsedParams = parseRpcRequestParams(request, method, envelopeMeta) if (parsedParams.error) { - reply(JSON.stringify(parsedParams.error)) + reply(JSON.stringify(boundLinearListReply(request, parsedParams.error))) return } @@ -107,6 +113,7 @@ export class RpcStreamingDispatcher { return method.handler(effectiveParams, { runtime, signal: options?.signal, + retainUntilDelivery: options?.retainUntilDelivery, requestId: request.id, connectionId: options?.connectionId, clientId: options?.clientId, @@ -140,9 +147,17 @@ export class RpcStreamingDispatcher { legacyCoordinator?.mutationCallerFingerprint ?? authenticatedCallerFingerprint ) recordRuntimeFeatureInteraction(runtime, request.method, result, undefined, request.params) - reply(JSON.stringify(successResponse(request.id, envelopeMeta, result))) + reply( + JSON.stringify( + boundLinearListReply(request, successResponse(request.id, envelopeMeta, result)) + ) + ) } catch (error) { - reply(JSON.stringify(mapDispatcherError(request, envelopeMeta, error))) + reply( + JSON.stringify( + boundLinearListReply(request, mapDispatcherError(request, envelopeMeta, error)) + ) + ) } return } @@ -160,6 +175,7 @@ export class RpcStreamingDispatcher { { runtime, signal: options?.signal, + retainUntilDelivery: options?.retainUntilDelivery, requestId: request.id, connectionId: options?.connectionId, clientId: options?.clientId, @@ -183,7 +199,11 @@ export class RpcStreamingDispatcher { request.params ) } catch (error) { - reply(JSON.stringify(mapDispatcherError(request, envelopeMeta, error))) + reply( + JSON.stringify( + boundLinearListReply(request, mapDispatcherError(request, envelopeMeta, error)) + ) + ) } } } diff --git a/src/main/runtime/rpc/transport.ts b/src/main/runtime/rpc/transport.ts index 8eedd20c0fb..8181907d7a9 100644 --- a/src/main/runtime/rpc/transport.ts +++ b/src/main/runtime/rpc/transport.ts @@ -12,6 +12,7 @@ // out. `startKeepalive` is opt-in per request — only long-poll dispatches // call it, so short RPCs pay no timer overhead. See design doc §3.1. export type RpcMessageContext = { + retainUntilDelivery?: (release: () => void) => void signal: AbortSignal startKeepalive: () => void } diff --git a/src/main/runtime/rpc/unix-socket-transport.ts b/src/main/runtime/rpc/unix-socket-transport.ts index ed66f94d907..f27dff1ba87 100644 --- a/src/main/runtime/rpc/unix-socket-transport.ts +++ b/src/main/runtime/rpc/unix-socket-transport.ts @@ -162,6 +162,17 @@ export class UnixSocketTransport implements RpcTransport { // handlers (e.g. orchestration.check --wait) arm it. See §3.1. private dispatchMessage(socket: Socket, rawMessage: string, inflight: Set<() => void>): void { let replied = false + let delivered = false + const retained = new Set<() => void>() + const settleDelivery = (): void => { + delivered = true + for (const release of retained) { + release() + } + retained.clear() + socket.off('close', settleDelivery) + } + socket.once('close', settleDelivery) let keepaliveTimer: NodeJS.Timeout | null = null // Why: each dispatch needs its own abort signal and keepalive timer // cleanup. Socket close runs every cleanup without touching sibling @@ -192,7 +203,9 @@ export class UnixSocketTransport implements RpcTransport { replied = true cleanupDispatch(false) if (!socket.destroyed && socket.writable) { - socket.write(`${response}\n`) + socket.write(`${response}\n`, settleDelivery) + } else { + settleDelivery() } } @@ -215,6 +228,13 @@ export class UnixSocketTransport implements RpcTransport { this.messageHandler?.(rawMessage, reply, { signal: abortController.signal, + retainUntilDelivery: (release) => { + if (delivered) { + release() + } else { + retained.add(release) + } + }, startKeepalive }) } diff --git a/src/main/runtime/runtime-linear-read-commands.ts b/src/main/runtime/runtime-linear-read-commands.ts index 4ef2f7d27c1..dcb13344219 100644 --- a/src/main/runtime/runtime-linear-read-commands.ts +++ b/src/main/runtime/runtime-linear-read-commands.ts @@ -196,9 +196,12 @@ export class RuntimeLinearReadCommands extends RuntimeLinearContextCommands { } } - async linearMcpIssueList(params: LinearMcpIssueListRequest): Promise { + async linearMcpIssueList( + params: LinearMcpIssueListRequest, + options: { signal?: AbortSignal; retainUntilDelivery?: (release: () => void) => void } = {} + ): Promise { try { - return await listMcpIssues(params) + return await listMcpIssues(params, options) } catch (error) { throw this.mapLinearReadFailure(error) } diff --git a/src/main/runtime/runtime-rpc/runtime-rpc-lifecycle.ts b/src/main/runtime/runtime-rpc/runtime-rpc-lifecycle.ts index 9cfe31251b6..41f633f523e 100644 --- a/src/main/runtime/runtime-rpc/runtime-rpc-lifecycle.ts +++ b/src/main/runtime/runtime-rpc/runtime-rpc-lifecycle.ts @@ -1,3 +1,5 @@ +import { isLinearPageRequest, linearListDeliveryFailure } from '../rpc/linear-list-reply-budget' +import type { RpcRequest } from '../rpc/core' import type { RuntimeTransportMetadata } from '../../../shared/runtime-bootstrap' import { watchRuntimeMetadataOwnership } from '../runtime-metadata-ownership-watch' import type { RpcTransport } from '../rpc/transport' @@ -50,6 +52,15 @@ export class RuntimeRpcLifecycle extends RuntimeRpcWebSocketDispatch { reply(JSON.stringify(response)) }) .catch((error) => { + try { + const request = JSON.parse(msg) as RpcRequest + if (isLinearPageRequest(request)) { + reply(JSON.stringify(linearListDeliveryFailure(request))) + return + } + } catch { + /* Invalid requests use the ordinary correlation fallback. */ + } const message = error instanceof Error ? error.message : String(error) // Why: best-effort id recovery so the client can correlate the error frame to its pending request. let id = 'unknown' diff --git a/src/main/runtime/runtime-rpc/runtime-rpc-request-admission.ts b/src/main/runtime/runtime-rpc/runtime-rpc-request-admission.ts index b79b87a5be5..1c6173ea2d3 100644 --- a/src/main/runtime/runtime-rpc/runtime-rpc-request-admission.ts +++ b/src/main/runtime/runtime-rpc/runtime-rpc-request-admission.ts @@ -1,3 +1,4 @@ +import { isLinearPageRequest } from '../rpc/linear-list-reply-budget' import type { RuntimeMetadata } from '../../../shared/runtime-bootstrap' import { writeRuntimeMetadata } from '../runtime-metadata' import type { RpcMessageContext } from '../rpc/transport' @@ -36,7 +37,8 @@ export class RuntimeRpcRequestAdmission extends RuntimeRpcBinaryRouting { try { return await this.dispatcher.dispatch(request, { - signal: longPoll ? context?.signal : undefined + signal: longPoll || isLinearPageRequest(request) ? context?.signal : undefined, + retainUntilDelivery: context?.retainUntilDelivery }) } finally { this.releaseLongPoll(longPoll) diff --git a/src/main/runtime/runtime-rpc/runtime-rpc-websocket-dispatch.ts b/src/main/runtime/runtime-rpc/runtime-rpc-websocket-dispatch.ts index dd714b77528..52394541303 100644 --- a/src/main/runtime/runtime-rpc/runtime-rpc-websocket-dispatch.ts +++ b/src/main/runtime/runtime-rpc/runtime-rpc-websocket-dispatch.ts @@ -101,10 +101,24 @@ export class RuntimeRpcWebSocketDispatch extends RuntimeRpcRequestAdmission { const abortRegistration = ws ? this.registerWebSocketDispatchAbort(ws) : null // Why: older pairings may lack scope metadata, so stamp the authenticated scope onto status.get. + const retained = new Set<() => void>() + const settleDelivery = (): void => { + for (const release of retained) { + release() + } + retained.clear() + } + const replyWithTransfer = (response: string): void => { + try { + reply(response) + } finally { + settleDelivery() + } + } const replyForRequest = request.method === 'status.get' - ? (response: string): void => reply(injectDeviceScope(response, device.scope)) - : reply + ? (response: string): void => replyWithTransfer(injectDeviceScope(response, device.scope)) + : replyWithTransfer const connectionId = ws ? this.mobileSocketWiring?.getConnectionId(ws) : undefined const pairingProvider = this.mobileRelayPairingProvider @@ -149,6 +163,9 @@ export class RuntimeRpcWebSocketDispatch extends RuntimeRpcRequestAdmission { : undefined, pairing: pairingContext, signal: abortRegistration?.signal, + retainUntilDelivery: (release) => { + retained.add(release) + }, sendBinary, registerBinaryStreamHandler: (streamId, handler) => this.registerBinaryStreamHandler(connectionId, streamId, handler), @@ -156,6 +173,7 @@ export class RuntimeRpcWebSocketDispatch extends RuntimeRpcRequestAdmission { this.registerBinaryMessageHandler(connectionId, handler) }) } finally { + settleDelivery() abortRegistration?.dispose() this.releaseLongPoll(longPoll, device.deviceId) } diff --git a/src/main/ssh/linear-list-ssh-delivery.ts b/src/main/ssh/linear-list-ssh-delivery.ts new file mode 100644 index 00000000000..1fc85c38762 --- /dev/null +++ b/src/main/ssh/linear-list-ssh-delivery.ts @@ -0,0 +1,96 @@ +import { stringifyJsonWithinByteLimit } from '../../shared/node-bounded-json-stringify' +import { linearListDeliveryFailure } from '../runtime/rpc/linear-list-reply-budget' +import { parseRemoteCliArgs } from './ssh-remote-cli-args' +import type { JsonRpcRequest, JsonRpcResponse } from './relay-protocol' +import type { RpcRequest } from '../runtime/rpc/core' + +export type LinearListDeliveryContext = { + signal: AbortSignal + retainUntilDelivery: (release: () => void) => void +} + +export class LinearListSshDelivery implements LinearListDeliveryContext { + private readonly controller = new AbortController() + private readonly retained = new Set<() => void>() + private settled = false + readonly signal = this.controller.signal + private constructor( + private readonly input: RpcRequest, + private readonly json: boolean + ) {} + + static forRequest(request: JsonRpcRequest): LinearListSshDelivery | undefined { + if (request.method !== 'orca.cli' || !Array.isArray(request.params?.argv)) { + return undefined + } + const argv = request.params.argv + if (!argv.every((value): value is string => typeof value === 'string')) { + return undefined + } + const parsed = parseRemoteCliArgs(argv) + if (parsed.commandPath.join(' ') !== 'linear list-issues' || parsed.flags.has('help')) { + return undefined + } + const continuation = parsed.flags.get('page-recovery') + return new LinearListSshDelivery( + { + id: 'remote-cli', + authToken: '', + method: 'linear.mcpListIssues', + params: { + workspaceId: parsed.flags.get('workspace'), + cursor: parsed.flags.get('cursor'), + ...(typeof continuation === 'string' + ? { pageRecovery: { version: 1, continuation } } + : {}) + } + }, + parsed.flags.has('json') + ) + } + + readonly retainUntilDelivery = (release: () => void): void => { + if (this.settled) { + release() + } else { + this.retained.add(release) + } + } + + readonly finish = (): void => { + this.settled = true + for (const release of this.retained) { + release() + } + this.retained.clear() + } + + readonly abort = (): void => { + this.controller.abort() + this.finish() + } + + bound(response: JsonRpcResponse): JsonRpcResponse { + try { + if (!Number.isSafeInteger(response.id)) { + throw new Error('invalid correlation') + } + stringifyJsonWithinByteLimit(response, 2_129_919, 2) + stringifyJsonWithinByteLimit(response, 2_129_920) + return response + } catch { + const failure = linearListDeliveryFailure(this.input) + return { + jsonrpc: '2.0', + id: Number.isSafeInteger(response.id) ? response.id : 0, + result: { + stdout: this.json ? `${JSON.stringify(failure, null, 2)}\n` : '', + stderr: this.json + ? '' + : 'Linear reply could not be delivered; retry the input position or restart concrete workspaces and reconcile by issue ID.\n', + exitCode: 1 + } + } + } + } +} diff --git a/src/main/ssh/ssh-channel-multiplexer.ts b/src/main/ssh/ssh-channel-multiplexer.ts index 3c9a8bd9ea1..c05b8b5ae5d 100644 --- a/src/main/ssh/ssh-channel-multiplexer.ts +++ b/src/main/ssh/ssh-channel-multiplexer.ts @@ -1,5 +1,6 @@ /* eslint-disable max-lines -- Why: the SSH relay protocol state machine keeps request, notification, keepalive, and cancellation semantics paired. */ +import { LinearListSshDelivery, type LinearListDeliveryContext } from './linear-list-ssh-delivery' import { FrameDecoder, MessageType, @@ -39,7 +40,10 @@ export type SshMultiplexerRequestOptions = { export type NotificationHandler = (method: string, params: Record) => void export type MethodNotificationHandler = (params: Record) => void -export type RequestHandler = (params: Record) => unknown +export type RequestHandler = ( + params: Record, + delivery?: LinearListDeliveryContext +) => unknown export type MultiplexerDisposeReason = 'shutdown' | 'connection_lost' @@ -499,15 +503,29 @@ export class SshChannelMultiplexer { return } + const delivery = LinearListSshDelivery.forRequest(msg) + const unsubscribe = delivery ? this.onDispose(delivery.abort) : undefined + const finish = (): void => { + unsubscribe?.() + delivery?.finish() + } + const send = (response: JsonRpcResponse): void => { + try { + this.sendMessage(delivery ? delivery.bound(response) : response, finish) + } catch (error) { + finish() + throw error + } + } try { - const result = await handler(msg.params ?? {}) - this.sendMessage({ + const result = await handler(msg.params ?? {}, delivery) + send({ jsonrpc: '2.0', id: msg.id, result: result ?? null }) } catch (err) { - this.sendMessage({ + send({ jsonrpc: '2.0', id: msg.id, error: { diff --git a/src/main/ssh/ssh-relay-session.ts b/src/main/ssh/ssh-relay-session.ts index a4fdb0f0fa7..3558ba34890 100644 --- a/src/main/ssh/ssh-relay-session.ts +++ b/src/main/ssh/ssh-relay-session.ts @@ -1433,7 +1433,7 @@ export class SshRelaySession { } private wireUpRemoteOrcaCli(mux: SshChannelMultiplexer, connectionIncarnation: string): void { - mux.onRequest('orca.cli', async (params) => { + mux.onRequest('orca.cli', async (params, delivery) => { if (!this.runtime) { throw new Error('Orca runtime is unavailable') } @@ -1465,7 +1465,8 @@ export class SshRelaySession { env, ...(stdin !== undefined ? { stdin } : {}), ...(artifactInput ? { artifactInput } : {}), - runtimeAuthority + runtimeAuthority, + delivery }) } finally { this.activeCompatibilityAttachmentIds.delete(runtimeAuthority.attachmentId) @@ -1495,7 +1496,8 @@ export class SshRelaySession { await acknowledgeRemoteOrcaCliPostOutput(this.runtime, { postOutput: parseRemoteOrcaCliPostOutput(params.postOutput), env, - runtimeAuthority + runtimeAuthority, + delivery }) return { acknowledged: true } } finally { diff --git a/src/main/ssh/ssh-remote-cli-host-passthrough.ts b/src/main/ssh/ssh-remote-cli-host-passthrough.ts index d62a23d8e30..bc5cd50de4e 100644 --- a/src/main/ssh/ssh-remote-cli-host-passthrough.ts +++ b/src/main/ssh/ssh-remote-cli-host-passthrough.ts @@ -1,3 +1,4 @@ +import type { LinearListDeliveryContext } from './linear-list-ssh-delivery' // The SSH shim runs the bundled CLI so remote shells get the full command surface. import { app } from 'electron' import { spawn as nodeSpawn } from 'node:child_process' @@ -27,6 +28,7 @@ export type SshCliRuntimeAuthority = { } export type RemoteOrcaCliRequest = { + delivery?: LinearListDeliveryContext argv: string[] cwd: string env: Record diff --git a/src/main/ssh/ssh-remote-cli-interactive-commands.ts b/src/main/ssh/ssh-remote-cli-interactive-commands.ts new file mode 100644 index 00000000000..5632295069d --- /dev/null +++ b/src/main/ssh/ssh-remote-cli-interactive-commands.ts @@ -0,0 +1,13 @@ +// Why: these commands run a foreground/interactive process attached to the +// caller's TTY (or a local tmux pane), which a buffered one-shot relay bridge +// cannot host. Everything else routes through the full host CLI. +export const HOST_INTERACTIVE_COMMANDS: Record = { + serve: + 'orca serve starts a foreground headless Orca server and cannot run through the SSH relay bridge. Run it directly on the machine that should host Orca.', + 'claude-teams': + 'orca claude-teams starts an interactive Claude Code session and cannot run through the SSH relay bridge. Run it in a terminal on the Orca host machine.', + 'agent-teams-tmux': + 'orca agent-teams-tmux is a tmux pane shim for the Orca host machine and cannot run through the SSH relay bridge.', + 'account add': + 'orca account add runs an interactive agent login and cannot run through the buffered SSH relay bridge. Run it directly in a terminal on the Orca host machine.' +} diff --git a/src/main/ssh/ssh-remote-linear-list-issues.ts b/src/main/ssh/ssh-remote-linear-list-issues.ts index 98f10e8f36c..7a847f75745 100644 --- a/src/main/ssh/ssh-remote-linear-list-issues.ts +++ b/src/main/ssh/ssh-remote-linear-list-issues.ts @@ -1,3 +1,4 @@ +import type { LinearConnectionStatus } from '../../shared/linear/workspace-types' import type { LinearMcpIssueListRequest } from '../../shared/linear/agent-access' import type { RpcDispatcher } from '../runtime/rpc/dispatcher' import type { RpcResponse } from '../runtime/rpc/core' @@ -41,6 +42,29 @@ export async function dispatchRemoteLinearListIssues( includeArchived: parsed.flags.get('include-archived') === true, workspaceId: optionalString(parsed.flags, 'workspace') } + const continuation = optionalString(parsed.flags, 'page-recovery') + if (continuation && (request.workspaceId !== 'all' || request.cursor)) { + throw new RemoteCliArgumentError( + 'invalid_argument', + '--page-recovery requires --workspace all and cannot use --cursor' + ) + } + if (request.workspaceId === 'all') { + const status = await dispatcher.dispatch({ + id: 'linear-page-capability', + authToken: 'remote-cli', + method: 'linear.status', + params: {} + }) + if (status.ok && (status.result as LinearConnectionStatus).mcpListPageRecoveryVersion === 1) { + request.pageRecovery = { version: 1, ...(continuation ? { continuation } : {}) } + } else if (continuation) { + throw new RemoteCliArgumentError( + 'linear_list_concrete_workspace_required', + 'This runtime does not support page recovery; restart concrete workspaces and reconcile by issue ID.' + ) + } + } return await dispatcher.dispatch({ id: `remote-cli-${Date.now()}`, authToken: 'remote-cli', diff --git a/src/main/ssh/ssh-remote-linear-output.ts b/src/main/ssh/ssh-remote-linear-output.ts index 533e743d10e..b9d419ca12a 100644 --- a/src/main/ssh/ssh-remote-linear-output.ts +++ b/src/main/ssh/ssh-remote-linear-output.ts @@ -284,6 +284,9 @@ function linearListWarnings( } function linearMcpListWarnings(result: LinearMcpIssueListResult): string { + if (result.meta.hasMore && result.meta.pageRecovery) { + return `warning: admitted batch; continue with --workspace all --page-recovery ${result.meta.pageRecovery.continuation}\n` + } const warnings = result.meta.workspaceErrors.map( (error) => `warning: ${error.workspace.name} unavailable for Linear: ${error.message}` ) diff --git a/src/main/ssh/ssh-remote-linear-read-flags.ts b/src/main/ssh/ssh-remote-linear-read-flags.ts index 08eea74c43d..3d4696b319b 100644 --- a/src/main/ssh/ssh-remote-linear-read-flags.ts +++ b/src/main/ssh/ssh-remote-linear-read-flags.ts @@ -69,6 +69,7 @@ export const LINEAR_MCP_ISSUE_LIST_FLAGS = new Set([ 'query', 'state', 'cursor', + 'page-recovery', 'order-by', 'project', 'release', diff --git a/src/main/ssh/ssh-remote-orca-cli.ts b/src/main/ssh/ssh-remote-orca-cli.ts index dba1bc1e2d5..d706e581afa 100644 --- a/src/main/ssh/ssh-remote-orca-cli.ts +++ b/src/main/ssh/ssh-remote-orca-cli.ts @@ -1,3 +1,4 @@ +import { HOST_INTERACTIVE_COMMANDS } from './ssh-remote-cli-interactive-commands' import type { CliStatusResult, RuntimeStatus } from '../../shared/runtime-types' import { runtimeHostConnectionState } from '../../shared/runtime-host-connection-state' import { projectRemoteAppStatus } from '../../shared/cli-app-status-projection' @@ -35,20 +36,6 @@ import { formatInProcessRemoteCliResult } from './ssh-remote-cli-in-process-resu export type { RemoteOrcaCliRequest, RemoteOrcaCliResult } from './ssh-remote-cli-host-passthrough' -// Why: these commands run a foreground/interactive process attached to the -// caller's TTY (or a local tmux pane), which a buffered one-shot relay bridge -// cannot host. Everything else routes through the full host CLI. -const HOST_INTERACTIVE_COMMANDS: Record = { - serve: - 'orca serve starts a foreground headless Orca server and cannot run through the SSH relay bridge. Run it directly on the machine that should host Orca.', - 'claude-teams': - 'orca claude-teams starts an interactive Claude Code session and cannot run through the SSH relay bridge. Run it in a terminal on the Orca host machine.', - 'agent-teams-tmux': - 'orca agent-teams-tmux is a tmux pane shim for the Orca host machine and cannot run through the SSH relay bridge.', - 'account add': - 'orca account add runs an interactive agent login and cannot run through the buffered SSH relay bridge. Run it directly in a terminal on the Orca host machine.' -} - export async function runRemoteOrcaCli( runtime: OrcaRuntimeService, request: RemoteOrcaCliRequest, @@ -82,6 +69,16 @@ export async function runRemoteOrcaCli( ) } + if (command === 'linear list-issues' && !parsed.flags.has('help')) { + request.delivery?.signal.throwIfAborted() + return runLegacyRemoteOrcaCli( + runtime, + request, + parsed, + json, + new HostCliUnavailableError('Linear page delivery is owned by the runtime dispatcher') + ) + } let passthroughFailure: HostCliUnavailableError | null = null try { return await runHostOrcaCliPassthrough(request, passthroughOptions) @@ -104,7 +101,11 @@ async function runLegacyRemoteOrcaCli( json: boolean, passthroughFailure: HostCliUnavailableError ): Promise { - const dispatcher = new RpcDispatcher({ runtime, methods: ALL_RPC_METHODS }) + const dispatcher = new RpcDispatcher({ + runtime, + methods: ALL_RPC_METHODS, + linearListDelivery: request.delivery + }) const help = getRemoteLinearHelp(parsed) if (help) { return { stdout: `${help}\n`, stderr: '', exitCode: 0 } diff --git a/src/shared/fetch-response-body.ts b/src/shared/fetch-response-body.ts index e9467fd3a2d..b4d15ba0f55 100644 --- a/src/shared/fetch-response-body.ts +++ b/src/shared/fetch-response-body.ts @@ -46,7 +46,8 @@ async function cancelReader(reader: ReadableStreamDefaultReader): Pr export async function readFetchResponseBytesWithinLimit( response: Response, - maxBytes = API_RESPONSE_MAX_BYTES + maxBytes = API_RESPONSE_MAX_BYTES, + signal?: AbortSignal ): Promise { if (!Number.isSafeInteger(maxBytes) || maxBytes < 0) { throw new RangeError('Response body limit must be a non-negative safe integer') @@ -64,9 +65,20 @@ export async function readFetchResponseBytesWithinLimit( const reader = response.body.getReader() let output = new Uint8Array(Math.min(maxBytes, INITIAL_RESPONSE_CAPACITY_BYTES)) let byteLength = 0 + let cancellation: Promise | undefined + const abort = (): void => { + output = new Uint8Array() + cancellation ??= cancelReader(reader) + } + signal?.addEventListener('abort', abort, { once: true }) + if (signal?.aborted) { + abort() + } try { while (true) { + signal?.throwIfAborted() const { done, value } = await reader.read() + signal?.throwIfAborted() if (done) { return output.subarray(0, byteLength) } @@ -88,6 +100,8 @@ export async function readFetchResponseBytesWithinLimit( byteLength = nextLength } } finally { + signal?.removeEventListener('abort', abort) + await cancellation reader.releaseLock() } } diff --git a/src/shared/linear/agent-access.ts b/src/shared/linear/agent-access.ts index 5e69edc0c55..1e2af4ddbcd 100644 --- a/src/shared/linear/agent-access.ts +++ b/src/shared/linear/agent-access.ts @@ -34,7 +34,16 @@ export const LINEAR_ERROR_CODES = [ 'linear_permission_denied', 'linear_auth_expired', 'linear_network_error', - 'linear_partial' + 'linear_partial', + 'linear_list_record_too_large', + 'linear_list_acquisition_too_large', + 'linear_list_capacity', + 'linear_list_metadata_capacity', + 'linear_list_invalid_response', + 'linear_list_cursor_cycle', + 'linear_list_empty_page', + 'linear_list_stale_recovery', + 'linear_list_concrete_workspace_required' ] as const export type LinearErrorCode = (typeof LINEAR_ERROR_CODES)[number] diff --git a/src/shared/linear/mcp-issue-list.ts b/src/shared/linear/mcp-issue-list.ts index 2419c7c2335..33683cd76bf 100644 --- a/src/shared/linear/mcp-issue-list.ts +++ b/src/shared/linear/mcp-issue-list.ts @@ -8,6 +8,7 @@ export type LinearMcpIssueListRequest = { limit?: number query?: string state?: string + pageRecovery?: { version: 1; continuation?: string } cursor?: string orderBy?: 'createdAt' | 'updatedAt' project?: string @@ -33,6 +34,14 @@ export type LinearMcpIssueListResult = { returned: number hasMore: boolean nextCursor?: string + pageRecovery?: { + version: 1 + continuation: string + ordering: 'admitted_batch' + consistency: 'best_effort' + stopReason?: string + } + omittedWorkspaceErrors?: number orderBy: 'createdAt' | 'updatedAt' workspaceId?: (string & {}) | 'all' partial: boolean @@ -40,6 +49,7 @@ export type LinearMcpIssueListResult = { workspace: LinearWorkspaceCandidate code: LinearErrorCode message: string + data?: unknown }[] } } diff --git a/src/shared/linear/workspace-types.ts b/src/shared/linear/workspace-types.ts index 508fd980a1e..5196777c2ee 100644 --- a/src/shared/linear/workspace-types.ts +++ b/src/shared/linear/workspace-types.ts @@ -31,6 +31,7 @@ export type LinearCollectionResult = { } export type LinearConnectionStatus = { + mcpListPageRecoveryVersion?: 1 connected: boolean viewer: LinearViewer | null workspaces?: LinearWorkspace[]