Files
orca/src/shared/runtime-rpc-call-queue.ts
T
OrcaWinandm4air 5f308bfa9c revert: take the 26 Phase 3 (#16741 port) PRs back out of main (#24559)
* Revert "feat(orcad): source-side dormant export of a relay-hosted SSH target (#16741 T6-8) (#24519)"

This reverts commit 783101b304.

* Revert "feat(ssh): update, roll back, recover and stop a managed orcad server (#16741 T6-5 follow-up) (#24463)"

This reverts commit 38c2d1dcb9.

* Revert "feat(ssh): deploy and pair an empty managed orcad server over SSH (#16741 T6-5) (#24453)"

This reverts commit 8b76683b40.

* Revert "fix(ssh): orcad GC honors the activation journal; readiness requires proven daemon coverage (#16741 T6 follow-up) (#24451)"

This reverts commit d3f8c5063b.

* Revert "feat(ssh): remote orcad stop by request file and journaled decommission (#16741 T6-4) (#24449)"

This reverts commit 43d9b43d3f.

* Revert "feat(orcad): supervisable server: stop requests, managed stop receipts and a lifetime that keeps its lock on failed teardown (#16741 T6-3) (#24433)"

This reverts commit b093d3ab20.

* Revert "feat(ssh): crash-safe orcad activation, rollback and recovery (#16741 T6-2) (#24423)"

This reverts commit 1a9ac0e955.

* Revert "feat(runtime): SSH access links for paired servers in a downgrade-safe sidecar (#16741 T5-1+T5-2) (#24420)"

This reverts commit 99db2bfae4.

* Revert "feat(relay): capability-gated owner reset with a durable preparation journal (#16741 T3 R1) (#24418)"

This reverts commit 34a582bd39.

* Revert "feat(ssh): track connection-manager drains, test probes and provider continuations (#16741 T2 P3+P8a) (#24407)"

This reverts commit d53063d2b1.

* Revert "feat(daemon): idle retirement, session census and recovery-only provider (#16741 T2 P4b) (#24409)"

This reverts commit ff212dbbef.

* Revert "feat(ssh): add pty.resumeClient and split SSH PTY process listing (#16741 T2 P5+P6) (#24414)"

This reverts commit 92cb71765e.

* Revert "feat(relay): await owned watcher and agent children on shutdown (#16741 T2 P1) (#24400)"

This reverts commit 6b36e4f85b.

* Revert "feat(session): retry failed renderer session writes and verify local folder PTYs (#16741 T2 P9) (#24406)"

This reverts commit d23ecef301.

* Revert "feat(ssh): remote orcad primitives on the pinned Node runtime (#16741 T6-1) (#24419)"

This reverts commit dd87ae578d.

* Revert "fix(runtime): fence runtime-environment subscriptions and status probes by identity (#16741 T5-3) (#24421)"

This reverts commit ece9e4d2e3.

* Revert "feat(orcad): migration manifest and dormant-state contracts (#16741 T6-7) (#24422)"

This reverts commit 3fbdaba262.

* Revert "feat(ssh): wire SshConnection through the work and transport close ledgers (#16741 T2 P2) (#24401)"

This reverts commit 4e8edc8872.

* Revert "feat(profiles): carry markdown frontmatter visibility in project transfers (#16741 T2 P7) (#24405)"

This reverts commit 60c93263cc.

* Revert "fix(runtime): project the PTY incarnation onto mobile session tabs (#24413)"

This reverts commit 99e0303572.

* Revert "feat(daemon): tag daemon stream data with the PTY incarnation id (#16741 T2 P4a) (#24402)"

This reverts commit 817af768b0.

* Revert "feat(ssh): port the SSH connection work ledger and transport close ledger (#16741 T2) (#24210)"

This reverts commit c9918931c8.

* Revert "feat(relay): fence and drain file and git response streams on shutdown (#24185)"

This reverts commit dc08ffeba9.

* Revert "refactor(runtime-rpc): extract the Node WebSocket lifecycle; opt-in pinned port (#24186)"

This reverts commit a789233bbb.

* Revert "feat(relay): route relay handlers through work admission; producer publication drain (#24181)"

This reverts commit 0b812bd698.

* Revert "feat(relay): land the #16741 T1 seam (work drain, publication drain, release gate) (#24156)"

This reverts commit 3aa2d3af7c.

---------

Co-authored-by: m4air <m4air@Mac.localdomain>
2026-10-02 00:52:32 -07:00

265 lines
8.8 KiB
TypeScript

import { abortSignalReason } from './abort-signal-reason'
import { REMOTE_RUNTIME_MAX_PREPARED_RPC_BYTES } from './remote-runtime-memory-limits'
const DEFAULT_REMOTE_RUNTIME_CALL_CONCURRENCY = 8
const DEFAULT_REMOTE_RUNTIME_BACKGROUND_CALL_CONCURRENCY = 2
export const RUNTIME_RPC_MAX_QUEUED_CALLS_PER_SELECTOR = 256
export const RUNTIME_RPC_MAX_QUEUED_CALLS_TOTAL = 2_048
export const RUNTIME_RPC_QUEUE_OVERLOAD_CODE = 'runtime_rpc_queue_overloaded'
export class RuntimeRpcCallQueueOverloadError extends Error {
readonly code = RUNTIME_RPC_QUEUE_OVERLOAD_CODE
constructor(readonly scope: 'selector' | 'global' | 'memory') {
super('Remote runtime call queue is full; retry after current calls finish.')
this.name = 'RuntimeRpcCallQueueOverloadError'
}
}
type QueuedRuntimeCall<T> = {
background: boolean
retainedBytes: number
run: () => Promise<T>
resolve: (value: T) => void
reject: (error: unknown) => void
started: boolean
signal?: AbortSignal
onAbort?: () => void
}
type RuntimeCallQueue = {
active: number
backgroundActive: number
foreground: (QueuedRuntimeCall<unknown> | undefined)[]
foregroundHead: number
background: (QueuedRuntimeCall<unknown> | undefined)[]
backgroundHead: number
}
export function isBackgroundRuntimeMethod(method: string): boolean {
return (
method === 'hostedReview.forBranch' ||
method === 'github.prForBranch' ||
method === 'github.listWorkItems' ||
method === 'github.countWorkItems' ||
method === 'git.status' ||
method === 'git.history' ||
method === 'git.conflictOperation' ||
method === 'git.branchCompare' ||
method === 'git.upstreamStatus' ||
method === 'worktree.prefetchCreateBase'
)
}
// Why its own lane: worktree.rm replies only when Git has deleted the checkout, so its calls would
// hold the foreground slots listing refreshes need; the background lane's slots belong to status.
function isLongWaitRuntimeMethod(method: string): boolean {
return method === 'worktree.rm'
}
export class RuntimeRpcCallQueuePool {
private readonly queues = new Map<string, RuntimeCallQueue>()
private queuedCallCount = 0
private retainedCallBytes = 0
constructor(
private readonly concurrency = DEFAULT_REMOTE_RUNTIME_CALL_CONCURRENCY,
private readonly backgroundConcurrency = DEFAULT_REMOTE_RUNTIME_BACKGROUND_CALL_CONCURRENCY,
private readonly maxQueuedPerSelector = RUNTIME_RPC_MAX_QUEUED_CALLS_PER_SELECTOR,
private readonly maxQueuedTotal = RUNTIME_RPC_MAX_QUEUED_CALLS_TOTAL,
private readonly maxRetainedBytes = REMOTE_RUNTIME_MAX_PREPARED_RPC_BYTES
) {}
enqueue<T>(
selector: string,
method: string,
run: () => Promise<T>,
retainedBytes = 0,
signal?: AbortSignal
): Promise<T> {
if (signal?.aborted) {
return Promise.reject(abortSignalReason(signal))
}
// Same concurrency bound, counted apart from the selector's other calls; global caps still apply.
const queueKey = isLongWaitRuntimeMethod(method) ? `${selector}\u0000long-wait` : selector
if (this.queuedCallCount >= this.maxQueuedTotal) {
return Promise.reject(new RuntimeRpcCallQueueOverloadError('global'))
}
const existingQueue = this.queues.get(queueKey)
if (existingQueue && this.queuedCount(existingQueue) >= this.maxQueuedPerSelector) {
return Promise.reject(new RuntimeRpcCallQueueOverloadError('selector'))
}
if (
!Number.isSafeInteger(retainedBytes) ||
retainedBytes < 0 ||
this.retainedCallBytes + retainedBytes > this.maxRetainedBytes
) {
return Promise.reject(new RuntimeRpcCallQueueOverloadError('memory'))
}
const queue = this.getQueue(queueKey)
return new Promise<T>((resolve, reject) => {
const call: QueuedRuntimeCall<T> = {
background: isBackgroundRuntimeMethod(method),
retainedBytes,
run,
resolve,
reject,
started: false,
signal
}
const targetQueue = call.background ? queue.background : queue.foreground
targetQueue.push(call as QueuedRuntimeCall<unknown>)
this.queuedCallCount += 1
this.retainedCallBytes += retainedBytes
if (signal) {
call.onAbort = () => this.cancelQueuedCall(queueKey, queue, call)
signal.addEventListener('abort', call.onAbort, { once: true })
}
this.pump(queueKey, queue)
})
}
private getQueue(selector: string): RuntimeCallQueue {
let queue = this.queues.get(selector)
if (!queue) {
queue = {
active: 0,
backgroundActive: 0,
foreground: [],
foregroundHead: 0,
background: [],
backgroundHead: 0
}
this.queues.set(selector, queue)
}
return queue
}
private pump(selector: string, queue: RuntimeCallQueue): void {
while (queue.active < this.concurrency) {
let call = this.takeForeground(queue)
if (!call && queue.backgroundActive < this.backgroundConcurrency) {
call = this.takeBackground(queue)
}
if (!call) {
break
}
call.started = true
if (call.signal && call.onAbort) {
call.signal.removeEventListener('abort', call.onAbort)
}
queue.active += 1
if (call.background) {
queue.backgroundActive += 1
}
// Why: runtime streams and worktree actions share transport capacity with
// per-card status refreshes, so decorative calls must not stampede it.
let runPromise: Promise<unknown>
try {
runPromise = call.run()
} catch (error) {
// Why: callers rely on queued work starting immediately, but sync
// validation errors must still flow through the cleanup path.
runPromise = Promise.reject(error)
}
void runPromise.then(call.resolve, call.reject).finally(() => {
this.retainedCallBytes = Math.max(0, this.retainedCallBytes - call.retainedBytes)
queue.active = Math.max(0, queue.active - 1)
if (call.background) {
queue.backgroundActive = Math.max(0, queue.backgroundActive - 1)
}
if (queue.active === 0 && this.isEmpty(queue)) {
this.queues.delete(selector)
return
}
this.pump(selector, queue)
})
}
}
private cancelQueuedCall<T>(
selector: string,
queue: RuntimeCallQueue,
call: QueuedRuntimeCall<T>
): void {
if (call.started || !call.signal) {
return
}
const targetQueue = call.background ? queue.background : queue.foreground
const head = call.background ? queue.backgroundHead : queue.foregroundHead
const index = targetQueue.indexOf(call as QueuedRuntimeCall<unknown>, head)
if (index === -1) {
return
}
targetQueue.splice(index, 1)
this.queuedCallCount = Math.max(0, this.queuedCallCount - 1)
this.retainedCallBytes = Math.max(0, this.retainedCallBytes - call.retainedBytes)
call.reject(abortSignalReason(call.signal))
if (queue.active === 0 && this.isEmpty(queue)) {
this.queues.delete(selector)
return
}
this.pump(selector, queue)
}
private takeForeground(queue: RuntimeCallQueue): QueuedRuntimeCall<unknown> | undefined {
if (queue.foregroundHead >= queue.foreground.length) {
return undefined
}
const call = queue.foreground[queue.foregroundHead]
queue.foreground[queue.foregroundHead] = undefined
queue.foregroundHead += 1
this.queuedCallCount = Math.max(0, this.queuedCallCount - 1)
this.compactForeground(queue)
return call
}
private takeBackground(queue: RuntimeCallQueue): QueuedRuntimeCall<unknown> | undefined {
if (queue.backgroundHead >= queue.background.length) {
return undefined
}
const call = queue.background[queue.backgroundHead]
queue.background[queue.backgroundHead] = undefined
queue.backgroundHead += 1
this.queuedCallCount = Math.max(0, this.queuedCallCount - 1)
this.compactBackground(queue)
return call
}
private compactForeground(queue: RuntimeCallQueue): void {
if (queue.foregroundHead <= 32 || queue.foregroundHead * 2 < queue.foreground.length) {
return
}
// Head indexes avoid repeated shifts; compaction bounds the consumed prefix.
queue.foreground.splice(0, queue.foregroundHead)
queue.foregroundHead = 0
}
private compactBackground(queue: RuntimeCallQueue): void {
if (queue.backgroundHead <= 32 || queue.backgroundHead * 2 < queue.background.length) {
return
}
queue.background.splice(0, queue.backgroundHead)
queue.backgroundHead = 0
}
private isEmpty(queue: RuntimeCallQueue): boolean {
return (
queue.foregroundHead >= queue.foreground.length &&
queue.backgroundHead >= queue.background.length
)
}
private queuedCount(queue: RuntimeCallQueue): number {
return (
queue.foreground.length -
queue.foregroundHead +
queue.background.length -
queue.backgroundHead
)
}
}