mirror of
https://github.com/stablyai/orca.git
synced 2026-10-07 16:02:29 +00:00
* Revert "feat(orcad): source-side dormant export of a relay-hosted SSH target (#16741 T6-8) (#24519)" This reverts commit783101b304. * Revert "feat(ssh): update, roll back, recover and stop a managed orcad server (#16741 T6-5 follow-up) (#24463)" This reverts commit38c2d1dcb9. * Revert "feat(ssh): deploy and pair an empty managed orcad server over SSH (#16741 T6-5) (#24453)" This reverts commit8b76683b40. * Revert "fix(ssh): orcad GC honors the activation journal; readiness requires proven daemon coverage (#16741 T6 follow-up) (#24451)" This reverts commitd3f8c5063b. * Revert "feat(ssh): remote orcad stop by request file and journaled decommission (#16741 T6-4) (#24449)" This reverts commit43d9b43d3f. * 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 commitb093d3ab20. * Revert "feat(ssh): crash-safe orcad activation, rollback and recovery (#16741 T6-2) (#24423)" This reverts commit1a9ac0e955. * Revert "feat(runtime): SSH access links for paired servers in a downgrade-safe sidecar (#16741 T5-1+T5-2) (#24420)" This reverts commit99db2bfae4. * Revert "feat(relay): capability-gated owner reset with a durable preparation journal (#16741 T3 R1) (#24418)" This reverts commit34a582bd39. * Revert "feat(ssh): track connection-manager drains, test probes and provider continuations (#16741 T2 P3+P8a) (#24407)" This reverts commitd53063d2b1. * Revert "feat(daemon): idle retirement, session census and recovery-only provider (#16741 T2 P4b) (#24409)" This reverts commitff212dbbef. * Revert "feat(ssh): add pty.resumeClient and split SSH PTY process listing (#16741 T2 P5+P6) (#24414)" This reverts commit92cb71765e. * Revert "feat(relay): await owned watcher and agent children on shutdown (#16741 T2 P1) (#24400)" This reverts commit6b36e4f85b. * Revert "feat(session): retry failed renderer session writes and verify local folder PTYs (#16741 T2 P9) (#24406)" This reverts commitd23ecef301. * Revert "feat(ssh): remote orcad primitives on the pinned Node runtime (#16741 T6-1) (#24419)" This reverts commitdd87ae578d. * Revert "fix(runtime): fence runtime-environment subscriptions and status probes by identity (#16741 T5-3) (#24421)" This reverts commitece9e4d2e3. * Revert "feat(orcad): migration manifest and dormant-state contracts (#16741 T6-7) (#24422)" This reverts commit3fbdaba262. * Revert "feat(ssh): wire SshConnection through the work and transport close ledgers (#16741 T2 P2) (#24401)" This reverts commit4e8edc8872. * Revert "feat(profiles): carry markdown frontmatter visibility in project transfers (#16741 T2 P7) (#24405)" This reverts commit60c93263cc. * Revert "fix(runtime): project the PTY incarnation onto mobile session tabs (#24413)" This reverts commit99e0303572. * Revert "feat(daemon): tag daemon stream data with the PTY incarnation id (#16741 T2 P4a) (#24402)" This reverts commit817af768b0. * Revert "feat(ssh): port the SSH connection work ledger and transport close ledger (#16741 T2) (#24210)" This reverts commitc9918931c8. * Revert "feat(relay): fence and drain file and git response streams on shutdown (#24185)" This reverts commitdc08ffeba9. * Revert "refactor(runtime-rpc): extract the Node WebSocket lifecycle; opt-in pinned port (#24186)" This reverts commita789233bbb. * Revert "feat(relay): route relay handlers through work admission; producer publication drain (#24181)" This reverts commit0b812bd698. * Revert "feat(relay): land the #16741 T1 seam (work drain, publication drain, release gate) (#24156)" This reverts commit3aa2d3af7c. --------- Co-authored-by: m4air <m4air@Mac.localdomain>
265 lines
8.8 KiB
TypeScript
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
|
|
)
|
|
}
|
|
}
|