diff --git a/src/relay/relay-work-admission.test.ts b/src/relay/relay-work-admission.test.ts new file mode 100644 index 00000000000..d5925953b11 --- /dev/null +++ b/src/relay/relay-work-admission.test.ts @@ -0,0 +1,223 @@ +import { describe, expect, it, vi } from 'vitest' +import type { RequestContext } from './dispatcher-contract' +import { RelayWorkAdmission } from './relay-work-admission' +import { SKILL_SSH_RELAY_CANCEL_UPLOAD_METHOD } from '../shared/skill-ssh-relay-contract' + +function context(): RequestContext { + return { clientId: 1, isStale: () => false } +} + +function pendingOperation() { + let resolve!: () => void + const promise = new Promise((settle) => { + resolve = settle + }) + return { promise, resolve } +} + +async function nextTurn(): Promise { + await new Promise((resolve) => setImmediate(resolve)) +} + +describe('RelayWorkAdmission', () => { + it('starts admitted operations synchronously and preserves their result', async () => { + const admission = new RelayWorkAdmission() + const operation = vi.fn(() => 'started') + const result = admission.run('fs.writeFile', context(), operation) + expect(operation).toHaveBeenCalledOnce() + await expect(result).resolves.toBe('started') + }) + + it('closes admission synchronously and drains existing work despite abort and disconnect', async () => { + const admission = new RelayWorkAdmission() + const controller = new AbortController() + let disconnected = false + const mutation = pendingOperation() + const running = admission.run( + 'fs.writeFile', + { ...context(), signal: controller.signal, isStale: () => disconnected }, + () => mutation.promise + ) + const drained = vi.fn() + const drain = admission.beginDrain().then(drained) + const rejectedOperation = vi.fn() + await expect(admission.run('fs.writeFile', context(), rejectedOperation)).rejects.toThrow( + 'relay_work_admission_closed' + ) + expect(rejectedOperation).not.toHaveBeenCalled() + controller.abort() + disconnected = true + await nextTurn() + expect(drained).not.toHaveBeenCalled() + mutation.resolve() + await running + await drain + expect(drained).toHaveBeenCalledOnce() + }) + + it.each([ + 'relay.status', + 'fs.unwatchAndWait', + 'agent.cancelExec', + SKILL_SSH_RELAY_CANCEL_UPLOAD_METHOD + ])('admits cleanup request %s while draining', async (method) => { + const admission = new RelayWorkAdmission() + const mutation = pendingOperation() + const running = admission.run('git.diff', context(), () => mutation.promise) + const drain = admission.beginDrain() + await expect( + admission.run(method, context(), () => { + mutation.resolve() + return 'acknowledged' + }) + ).resolves.toBe('acknowledged') + await running + await drain + }) + + it.each([ + 'rpc.cancel', + 'git.responseAck', + 'git.cancelResponseStream', + 'fs.streamAck', + 'fs.cancelStream', + 'fs.unwatch', + 'pty.ackData', + 'pty.setDeliveryPaused' + ])('admits cleanup notification %s while draining', async (method) => { + const admission = new RelayWorkAdmission() + const mutation = pendingOperation() + const running = admission.run('git.diff', context(), () => mutation.promise) + const drain = admission.beginDrain() + const acknowledge = vi.fn(() => mutation.resolve()) + admission.runNotification(method, context(), acknowledge) + expect(acknowledge).toHaveBeenCalledOnce() + await running + await drain + await expect(admission.run('git.diff', context(), () => {})).rejects.toThrow( + 'relay_work_admission_closed' + ) + }) + + it('excludes the exact reset request without excluding another request from the same client', async () => { + const admission = new RelayWorkAdmission() + const mutation = pendingOperation() + const running = admission.run('fs.writeFile', context(), () => mutation.promise) + const resetContext = context() + const drained = vi.fn() + const reset = admission.run('relay.reset', resetContext, async () => { + await admission.beginDrain(resetContext) + drained() + }) + await nextTurn() + expect(drained).not.toHaveBeenCalled() + mutation.resolve() + await running + await reset + expect(drained).toHaveBeenCalledOnce() + }) + + it('rejects a forged context without closing admission', async () => { + const admission = new RelayWorkAdmission() + const actualContext = context() + await admission.run('relay.reset', actualContext, async () => { + await expect(admission.beginDrain({ ...actualContext })).rejects.toThrow( + 'relay_work_drain_context_not_active' + ) + await expect(admission.run('fs.writeFile', context(), () => 'open')).resolves.toBe('open') + }) + }) + + it('rejects an already-settled request context without closing admission', async () => { + const admission = new RelayWorkAdmission() + const oldContext = context() + await admission.run('relay.reset', oldContext, () => {}) + await expect(admission.beginDrain(oldContext)).rejects.toThrow( + 'relay_work_drain_context_not_active' + ) + await expect(admission.run('fs.writeFile', context(), () => 'open')).resolves.toBe('open') + }) + + it('waits for reset work when the host lifecycle supplies no exclusion', async () => { + const admission = new RelayWorkAdmission() + const mutation = pendingOperation() + const reset = admission.run('relay.reset', context(), () => mutation.promise) + const drained = vi.fn() + const drain = admission.beginDrain().then(drained) + await nextTurn() + expect(drained).not.toHaveBeenCalled() + mutation.resolve() + await reset + await drain + expect(drained).toHaveBeenCalledOnce() + }) + + it('drops new mutation notifications without invoking their handlers during drain', async () => { + const admission = new RelayWorkAdmission() + await admission.beginDrain() + const mutate = vi.fn() + admission.runNotification('pty.write', context(), mutate) + expect(mutate).not.toHaveBeenCalled() + await expect(admission.run('pty.ackData', context(), mutate)).rejects.toThrow( + 'relay_work_admission_closed' + ) + }) + + it('tracks asynchronous notifications through settlement while draining', async () => { + const admission = new RelayWorkAdmission() + const mutation = pendingOperation() + const operation = vi.fn(async () => mutation.promise) + admission.runNotification('pty.write', context(), operation) + expect(operation).toHaveBeenCalledOnce() + const drained = vi.fn() + const drain = admission.beginDrain().then(drained) + await nextTurn() + expect(drained).not.toHaveBeenCalled() + mutation.resolve() + await drain + expect(drained).toHaveBeenCalledOnce() + }) + + it('releases tracking after synchronous notification failure', async () => { + const admission = new RelayWorkAdmission() + expect(() => + admission.runNotification('pty.write', context(), () => { + throw new Error('notification_failed') + }) + ).toThrow('notification_failed') + await admission.beginDrain() + const mutate = vi.fn() + admission.runNotification('pty.write', context(), mutate) + expect(mutate).not.toHaveBeenCalled() + }) + + it('settles concurrent drains after an admitted operation fails without reopening', async () => { + const admission = new RelayWorkAdmission() + const mutation = pendingOperation() + const running = admission.run('fs.writeFile', context(), async () => { + await mutation.promise + throw new Error('write_failed') + }) + const failure = expect(running).rejects.toThrow('write_failed') + const firstDrain = admission.beginDrain() + const secondDrain = admission.beginDrain() + mutation.resolve() + await failure + await Promise.all([firstDrain, secondDrain]) + await expect(admission.run('fs.writeFile', context(), () => {})).rejects.toThrow( + 'relay_work_admission_closed' + ) + }) + + it('turns synchronous operation failures into rejected promises and releases tracking', async () => { + const admission = new RelayWorkAdmission() + const result = admission.run('fs.writeFile', context(), () => { + throw new Error('write_failed') + }) + await expect(result).rejects.toThrow('write_failed') + await admission.beginDrain() + await expect(admission.run('fs.writeFile', context(), () => {})).rejects.toThrow( + 'relay_work_admission_closed' + ) + }) +}) diff --git a/src/relay/relay-work-admission.ts b/src/relay/relay-work-admission.ts new file mode 100644 index 00000000000..4003f755658 --- /dev/null +++ b/src/relay/relay-work-admission.ts @@ -0,0 +1,84 @@ +import type { RequestContext } from './dispatcher-contract' +import { allowsRelayWorkDuringDrain } from '../shared/relay-work-drain-contract' + +/** Handler drain is not proof of physical process exit or detached producer cleanup. */ +export class RelayWorkAdmission { + private draining = false + private readonly active = new Map }>() + + allows(method: string, notification = false): boolean { + return !this.draining || allowsRelayWorkDuringDrain(method, notification) + } + + async run( + method: string, + context: RequestContext, + operation: () => T | Promise + ): Promise { + if (!this.allows(method)) { + throw new Error('relay_work_admission_closed') + } + return this.track(context, operation) + } + + runNotification(method: string, context: RequestContext, operation: () => void): void { + if (!this.allows(method, true)) { + return + } + const finish = this.trackEntry(context) + try { + const result = operation() + void Promise.resolve(result).then(finish, (error) => { + finish() + process.stderr.write(`[relay] Notification handler failed: ${String(error)}\n`) + }) + } catch (error) { + finish() + throw error + } + } + + assertActiveContext(context: RequestContext): void { + if (![...this.active.values()].some((entry) => entry.context === context)) { + throw new Error('relay_work_drain_context_not_active') + } + } + + /** Only the actual initiating request context may be excluded from its own drain. */ + async beginDrain(exclude?: RequestContext): Promise { + if (exclude) { + this.assertActiveContext(exclude) + } + this.draining = true + for (;;) { + const pending = [...this.active.values()].filter((entry) => entry.context !== exclude) + if (pending.length === 0) { + return + } + await Promise.all(pending.map((entry) => entry.done)) + } + } + + private async track(context: RequestContext, operation: () => T | Promise): Promise { + const finish = this.trackEntry(context) + try { + return await operation() + } finally { + finish() + } + } + + // Why: no Promise.withResolvers — the legacy relay bundle still targets Node 18 hosts. + private trackEntry(context: RequestContext): () => void { + const key = {} + let resolve!: () => void + const done = new Promise((settle) => { + resolve = settle + }) + this.active.set(key, { context, done }) + return () => { + this.active.delete(key) + resolve() + } + } +} diff --git a/src/shared/pty-ownership-transfer-release-gate.test.ts b/src/shared/pty-ownership-transfer-release-gate.test.ts new file mode 100644 index 00000000000..dbfc79ed6ae --- /dev/null +++ b/src/shared/pty-ownership-transfer-release-gate.test.ts @@ -0,0 +1,29 @@ +import { describe, expect, it } from 'vitest' +import { + isPtyOwnershipTransferMutationEnabled, + PTY_OWNERSHIP_TRANSFER_CANARY_ENV +} from './pty-ownership-transfer-release-gate' + +describe('PTY ownership-transfer release gate', () => { + it('is closed unless the exact canary value is present', () => { + expect(isPtyOwnershipTransferMutationEnabled({})).toBe(false) + expect( + isPtyOwnershipTransferMutationEnabled({ + [PTY_OWNERSHIP_TRANSFER_CANARY_ENV]: 'true' + }) + ).toBe(false) + expect( + isPtyOwnershipTransferMutationEnabled({ + [PTY_OWNERSHIP_TRANSFER_CANARY_ENV]: '1 ' + }) + ).toBe(false) + }) + + it('opens only for an explicit canary opt-in', () => { + expect( + isPtyOwnershipTransferMutationEnabled({ + [PTY_OWNERSHIP_TRANSFER_CANARY_ENV]: '1' + }) + ).toBe(true) + }) +}) diff --git a/src/shared/pty-ownership-transfer-release-gate.ts b/src/shared/pty-ownership-transfer-release-gate.ts new file mode 100644 index 00000000000..2c411d89707 --- /dev/null +++ b/src/shared/pty-ownership-transfer-release-gate.ts @@ -0,0 +1,12 @@ +/** + * The live PTY ownership-transfer path is a release canary until routed, + * cross-platform recovery evidence is complete. Keep the opt-in exact and + * process-scoped so ordinary profiles cannot inherit it accidentally. + */ +export const PTY_OWNERSHIP_TRANSFER_CANARY_ENV = 'ORCA_ENABLE_PTY_OWNERSHIP_TRANSFER_MUTATION' + +export function isPtyOwnershipTransferMutationEnabled( + env: Record = process.env +): boolean { + return env[PTY_OWNERSHIP_TRANSFER_CANARY_ENV] === '1' +} diff --git a/src/shared/relay-work-drain-contract.ts b/src/shared/relay-work-drain-contract.ts new file mode 100644 index 00000000000..9b21e6bcfe5 --- /dev/null +++ b/src/shared/relay-work-drain-contract.ts @@ -0,0 +1,24 @@ +import { SKILL_SSH_RELAY_CANCEL_UPLOAD_METHOD } from './skill-ssh-relay-contract' + +// Why: a drain must still admit the requests and notifications that let in-flight work finish or cancel. +// Owner-reset (T3) and network-tunnel (T4) methods join these sets when their contracts land. +const drainRequests = new Set([ + 'relay.status', + 'fs.unwatchAndWait', + 'agent.cancelExec', + SKILL_SSH_RELAY_CANCEL_UPLOAD_METHOD +]) +const drainNotifications = new Set([ + 'rpc.cancel', + 'git.responseAck', + 'git.cancelResponseStream', + 'fs.streamAck', + 'fs.cancelStream', + 'fs.unwatch', + 'pty.ackData', + 'pty.setDeliveryPaused' +]) + +export function allowsRelayWorkDuringDrain(method: string, notification = false): boolean { + return (notification ? drainNotifications : drainRequests).has(method) +} diff --git a/src/shared/runtime-rpc-call-queue-retirement.test.ts b/src/shared/runtime-rpc-call-queue-retirement.test.ts new file mode 100644 index 00000000000..3958c200991 --- /dev/null +++ b/src/shared/runtime-rpc-call-queue-retirement.test.ts @@ -0,0 +1,65 @@ +import { expect, it, vi } from 'vitest' +import { RuntimeRpcCallQueueBusyError, RuntimeRpcCallQueuePool } from './runtime-rpc-call-queue' + +it('refuses an active or queued selector without interrupting or replaying its calls', async () => { + const queue = new RuntimeRpcCallQueuePool(1, 1) + const pending = Promise.withResolvers() + const first = queue.enqueue('host', 'terminal.send', () => pending.promise) + const secondRun = vi.fn(async () => 'second-ack') + const second = queue.enqueue('host', 'terminal.send', secondRun) + expect(() => queue.holdIdleSelectors(['host'])).toThrow(RuntimeRpcCallQueueBusyError) + pending.resolve('first-ack') + await expect(first).resolves.toBe('first-ack') + await expect(second).resolves.toBe('second-ack') + expect(secondRun).toHaveBeenCalledOnce() + await vi.waitFor(() => queue.holdIdleSelectors(['host'])()) +}) + +it('holds both identities, refuses new calls before dispatch, and leaves other hosts running', async () => { + const queue = new RuntimeRpcCallQueuePool() + const release = queue.holdIdleSelectors(['canonical', 'historical', 'canonical']) + const run = vi.fn(async () => 'ack') + for (const id of ['canonical', 'historical']) { + await expect(queue.enqueue(id, 'terminal.send', run)).rejects.toMatchObject({ + code: 'runtime_rpc_queue_busy' + }) + } + expect(run).not.toHaveBeenCalled() + await expect(queue.enqueue('other', 'terminal.send', run)).resolves.toBe('ack') + release() + await expect(queue.enqueue('historical', 'terminal.send', run)).resolves.toBe('ack') + expect(run).toHaveBeenCalledTimes(2) +}) + +it('acquires a group atomically without retaining a partial hold on failure', async () => { + const queue = new RuntimeRpcCallQueuePool() + const release = queue.holdIdleSelectors(['historical']) + expect(() => queue.holdIdleSelectors(['canonical', 'historical'])).toThrow( + RuntimeRpcCallQueueBusyError + ) + await expect(queue.enqueue('canonical', 'repo.list', async () => 'ok')).resolves.toBe('ok') + release() +}) + +it('does not let an old release clear a newer hold', async () => { + const queue = new RuntimeRpcCallQueuePool() + const first = queue.holdIdleSelectors(['host']) + first() + const second = queue.holdIdleSelectors(['host']) + first() + await expect(queue.enqueue('host', 'repo.list', async () => 'ok')).rejects.toThrow( + RuntimeRpcCallQueueBusyError + ) + second() + await expect(queue.enqueue('host', 'repo.list', async () => 'ok')).resolves.toBe('ok') +}) + +it('counts a selector busy while only its long-wait lane has a call', async () => { + const queue = new RuntimeRpcCallQueuePool() + const pending = Promise.withResolvers() + const removal = queue.enqueue('host', 'worktree.rm', () => pending.promise) + expect(() => queue.holdIdleSelectors(['host'])).toThrow(RuntimeRpcCallQueueBusyError) + pending.resolve('removed') + await expect(removal).resolves.toBe('removed') + await vi.waitFor(() => queue.holdIdleSelectors(['host'])()) +}) diff --git a/src/shared/runtime-rpc-call-queue.ts b/src/shared/runtime-rpc-call-queue.ts index 132c8099273..0503c5e06c0 100644 --- a/src/shared/runtime-rpc-call-queue.ts +++ b/src/shared/runtime-rpc-call-queue.ts @@ -16,6 +16,15 @@ export class RuntimeRpcCallQueueOverloadError extends Error { } } +export class RuntimeRpcCallQueueBusyError extends Error { + readonly code = 'runtime_rpc_queue_busy' + + constructor() { + super('Runtime calls are active or their routing is changing; retry after they settle.') + this.name = 'RuntimeRpcCallQueueBusyError' + } +} + type QueuedRuntimeCall = { background: boolean retainedBytes: number @@ -57,8 +66,13 @@ function isLongWaitRuntimeMethod(method: string): boolean { return method === 'worktree.rm' } +function longWaitQueueKey(selector: string): string { + return `${selector}\u0000long-wait` +} + export class RuntimeRpcCallQueuePool { private readonly queues = new Map() + private readonly heldSelectors = new Set() private queuedCallCount = 0 private retainedCallBytes = 0 @@ -70,6 +84,32 @@ export class RuntimeRpcCallQueuePool { private readonly maxRetainedBytes = REMOTE_RUNTIME_MAX_PREPARED_RPC_BYTES ) {} + /** Acquires all idle selectors atomically; release never replays refused calls. */ + holdIdleSelectors(selectors: readonly string[]): () => void { + const unique = [...new Set(selectors)] + // A selector's long-wait lane counts too: holding it must see every call still in flight. + const busy = (selector: string): boolean => + this.heldSelectors.has(selector) || + this.queues.has(selector) || + this.queues.has(longWaitQueueKey(selector)) + if (unique.some(busy)) { + throw new RuntimeRpcCallQueueBusyError() + } + for (const selector of unique) { + this.heldSelectors.add(selector) + } + let released = false + return () => { + if (released) { + return + } + released = true + for (const selector of unique) { + this.heldSelectors.delete(selector) + } + } + } + enqueue( selector: string, method: string, @@ -80,8 +120,11 @@ export class RuntimeRpcCallQueuePool { if (signal?.aborted) { return Promise.reject(abortSignalReason(signal)) } + if (this.heldSelectors.has(selector)) { + return Promise.reject(new RuntimeRpcCallQueueBusyError()) + } // Same concurrency bound, counted apart from the selector's other calls; global caps still apply. - const queueKey = isLongWaitRuntimeMethod(method) ? `${selector}\u0000long-wait` : selector + const queueKey = isLongWaitRuntimeMethod(method) ? longWaitQueueKey(selector) : selector if (this.queuedCallCount >= this.maxQueuedTotal) { return Promise.reject(new RuntimeRpcCallQueueOverloadError('global')) } diff --git a/src/shared/transport-publication-drain.test.ts b/src/shared/transport-publication-drain.test.ts new file mode 100644 index 00000000000..e515cf464f6 --- /dev/null +++ b/src/shared/transport-publication-drain.test.ts @@ -0,0 +1,70 @@ +import { describe, expect, it, vi } from 'vitest' +import { TransportPublicationDrain } from './transport-publication-drain' + +describe('TransportPublicationDrain', () => { + it('waits for every tracked write, including writes tracked while draining', async () => { + const publication = new TransportPublicationDrain(() => {}) + const first = publication.trackWrite() + const done = vi.fn() + const drain = publication.drain(new AbortController().signal).then(done) + const second = publication.trackWrite() + first({ ok: true }) + await Promise.resolve() + expect(done).not.toHaveBeenCalled() + expect(() => publication.assertDrained()).toThrow('transport_publication_not_drained') + second({ ok: true }) + await drain + expect(done).toHaveBeenCalledOnce() + }) + + it('ignores a repeated settlement of the same write', async () => { + const publication = new TransportPublicationDrain(() => {}) + const settle = publication.trackWrite() + const other = publication.trackWrite() + settle({ ok: true }) + settle({ ok: true }) + expect(() => publication.assertDrained()).toThrow('transport_publication_not_drained') + other({ ok: true }) + await publication.drain(new AbortController().signal) + }) + + it('aborts only the observer and keeps a later write failure sticky', async () => { + const onFailure = vi.fn() + const publication = new TransportPublicationDrain(() => {}, onFailure) + const settle = publication.trackWrite() + const controller = new AbortController() + const drain = publication.drain(controller.signal) + controller.abort(new Error('observer cancelled')) + await expect(drain).rejects.toThrow('observer cancelled') + settle({ ok: false, error: new Error('lost write') }) + expect(onFailure).toHaveBeenCalledOnce() + await expect(publication.drain(new AbortController().signal)).rejects.toThrow('lost write') + expect(() => publication.assertCurrent()).toThrow('lost write') + }) + + it('fails once when the transport it was bound to is replaced', async () => { + let current = true + const onFailure = vi.fn() + const publication = new TransportPublicationDrain(() => { + if (!current) { + throw new Error('transport_replaced') + } + }, onFailure) + const settle = publication.trackWrite() + current = false + const drain = publication.drain(new AbortController().signal) + await expect(drain).rejects.toThrow('transport_replaced') + settle({ ok: true }) + expect(() => publication.assertDrained()).toThrow('transport_replaced') + expect(onFailure).toHaveBeenCalledOnce() + }) + + it('refuses construction against a transport that is already stale', () => { + expect( + () => + new TransportPublicationDrain(() => { + throw new Error('transport_replaced') + }) + ).toThrow('transport_replaced') + }) +}) diff --git a/src/shared/transport-publication-drain.ts b/src/shared/transport-publication-drain.ts new file mode 100644 index 00000000000..3814b120195 --- /dev/null +++ b/src/shared/transport-publication-drain.ts @@ -0,0 +1,87 @@ +import { waitForPromiseWithSignal } from './abort-signal-reason' + +export type TransportPublicationSettlement = { ok: true } | { ok: false; error: Error } + +/** Local write completion only; downstream consumption requires separate proof. */ +export class TransportPublicationDrain { + private pending = 0 + protected failure: Error | undefined + private readonly observers = new Set<() => void>() + + constructor( + private readonly assertTransport: () => void, + private readonly onFailure: (error: Error) => void = () => {} + ) { + this.assertCurrent() + } + + trackWrite(): (result: TransportPublicationSettlement) => void { + this.pending++ + let settled = false + return (result) => { + if (settled) { + return + } + settled = true + this.pending-- + if (!result.ok) { + this.fail(result.error) + } + this.changed() + } + } + + fail(error: Error): void { + if (this.failure) { + return + } + this.failure = error + this.changed() + this.onFailure(error) + } + + assertDrained(): void { + this.assertCurrent() + if (this.pending !== 0) { + throw new Error('transport_publication_not_drained') + } + } + + async drain(signal: AbortSignal): Promise { + signal.throwIfAborted() + while (this.pending !== 0) { + this.assertCurrent() + // Why: no Promise.withResolvers — the legacy relay bundle still targets Node 18 hosts. + let notify!: () => void + const changed = new Promise((resolve) => { + notify = resolve + }) + this.observers.add(notify) + try { + await waitForPromiseWithSignal(changed, signal) + } finally { + this.observers.delete(notify) + } + } + signal.throwIfAborted() + this.assertDrained() + } + + assertCurrent(): void { + if (this.failure) { + throw this.failure + } + try { + this.assertTransport() + } catch (error) { + this.fail(error instanceof Error ? error : new Error(String(error))) + throw this.failure + } + } + + private changed(): void { + for (const notify of this.observers) { + notify() + } + } +}