Implement request-owned Linear page acquisition and recovery

This commit is contained in:
Merge Sim
2026-09-08 00:19:17 -07:00
parent 15adbd9d18
commit 0ff258490c
44 changed files with 1786 additions and 297 deletions
+19
View File
@@ -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<LinearConnectionStatus>('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<LinearMcpIssueListResult>('linear.mcpListIssues', request)
if (!json) {
printLinearMcpIssueListWarnings(response.result)
+5 -1
View File
@@ -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}`
+1
View File
@@ -53,6 +53,7 @@ export const LINEAR_MCP_COMMAND_SPECS: CommandSpec[] = [
'query',
'state',
'cursor',
'page-recovery',
'order-by',
'project',
'release',
+1
View File
@@ -179,6 +179,7 @@ export function getStatus(): LinearConnectionStatus {
return {
connected: state.workspaces.length > 0,
mcpListPageRecoveryVersion: 1,
viewer: activeWorkspace,
workspaces: state.workspaces,
activeWorkspaceId: state.activeWorkspaceId,
+7 -1
View File
@@ -181,12 +181,18 @@ export async function withLinearRead<T>(
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()
@@ -0,0 +1,29 @@
const reads = new Map<string, Set<AbortController>>()
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.'))
}
}
+44 -12
View File
@@ -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<void> {
export function acquire(signal?: AbortSignal): Promise<void> {
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
}
}
+3
View File
@@ -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))
@@ -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
@@ -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<string, unknown> & { first: number },
signal: AbortSignal
): Promise<z.infer<typeof envelope>['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
}
@@ -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<string>()
private readonly cursors = new Set<string>()
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')
}
@@ -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<T>() {
let resolve!: (value: T) => void
const promise = new Promise<T>((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<ReadableStreamReadResult<Uint8Array>>()
const cancel = deferred<void>()
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())
})
})
@@ -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<T>(workspaceId: string, operation: (signal: AbortSignal) => Promise<T>): Promise<T> {
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<never>((_, reject) => {
abort = () => reject(signal.reason)
signal.addEventListener('abort', abort, { once: true })
if (signal.aborted) {
abort()
}
})
])
} finally {
if (abort) {
signal.removeEventListener('abort', abort)
}
}
}
}
@@ -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<typeof row>[]) {
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<Response>((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'
})
})
})
+280
View File
@@ -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<LinearMcpIssueListResult> {
const admission = new IssueListAdmission()
const failures: LinearMcpIssueListResult['meta']['workspaceErrors'] = []
const failed = new Set<string>()
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 } : {})
}
}
}
+166
View File
@@ -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<typeof vector>
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
}
}
+31 -14
View File
@@ -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<typeof vi.fn> = 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()
+53 -215
View File
@@ -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<LinearMcpIssueListResult> {
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<typeof w> => !!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<WorkspaceIssuePage> {
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<string, unknown>
>(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] ?? '')
}
+1
View File
@@ -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
@@ -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
+19 -5
View File
@@ -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<RpcResponse> {
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))
}
}
@@ -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<RpcRequest, 'method' | 'params'>): 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)
}
}
@@ -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(
+17 -11
View File
@@ -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',
@@ -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<void> {
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))
)
)
}
}
}
+1
View File
@@ -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
}
+21 -1
View File
@@ -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
})
}
@@ -196,9 +196,12 @@ export class RuntimeLinearReadCommands extends RuntimeLinearContextCommands {
}
}
async linearMcpIssueList(params: LinearMcpIssueListRequest): Promise<LinearMcpIssueListResult> {
async linearMcpIssueList(
params: LinearMcpIssueListRequest,
options: { signal?: AbortSignal; retainUntilDelivery?: (release: () => void) => void } = {}
): Promise<LinearMcpIssueListResult> {
try {
return await listMcpIssues(params)
return await listMcpIssues(params, options)
} catch (error) {
throw this.mapLinearReadFailure(error)
}
@@ -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'
@@ -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)
@@ -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)
}
+96
View File
@@ -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
}
}
}
}
}
+22 -4
View File
@@ -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<string, unknown>) => void
export type MethodNotificationHandler = (params: Record<string, unknown>) => void
export type RequestHandler = (params: Record<string, unknown>) => unknown
export type RequestHandler = (
params: Record<string, unknown>,
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: {
+5 -3
View File
@@ -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 {
@@ -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<string, string>
@@ -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<string, string> = {
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.'
}
@@ -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',
+3
View File
@@ -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}`
)
@@ -69,6 +69,7 @@ export const LINEAR_MCP_ISSUE_LIST_FLAGS = new Set([
'query',
'state',
'cursor',
'page-recovery',
'order-by',
'project',
'release',
+16 -15
View File
@@ -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<string, string> = {
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<RemoteOrcaCliResult> {
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 }
+15 -1
View File
@@ -46,7 +46,8 @@ async function cancelReader(reader: ReadableStreamDefaultReader<Uint8Array>): Pr
export async function readFetchResponseBytesWithinLimit(
response: Response,
maxBytes = API_RESPONSE_MAX_BYTES
maxBytes = API_RESPONSE_MAX_BYTES,
signal?: AbortSignal
): Promise<Uint8Array> {
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<void> | 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()
}
}
+10 -1
View File
@@ -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]
+10
View File
@@ -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
}[]
}
}
+1
View File
@@ -31,6 +31,7 @@ export type LinearCollectionResult<T> = {
}
export type LinearConnectionStatus = {
mcpListPageRecoveryVersion?: 1
connected: boolean
viewer: LinearViewer | null
workspaces?: LinearWorkspace[]